1pub mod commands;
2pub mod history;
3pub mod state;
4pub mod undo;
5
6use std::path::{Path, PathBuf};
7use std::sync::Arc;
8use std::thread;
9
10use thiserror::Error;
11
12use crate::audio::{
13 analyzer::VizAnalyzer,
14 backend::{self, AudioBackend, AudioEngineHandle, BackendError, SampleRateWatch},
15 buffer, streaming,
16 viz::{VizBuffer, VizSnapshot},
17};
18use crate::remote::client::PlaybackReportState;
19use buffer::PlaybackTimeline;
20use commands::{CommandChannel, PlayerCommand};
21use history::{InFlight, PlayEvent, PlayRecorder, PlaybackReport};
22use state::{
23 ItemState, LoadState, PlaybackSource, PlaybackState, QueueItemId, SharedPlayerState, TrackInfo,
24};
25use undo::{UndoEntry, UndoStack};
26
27pub(crate) const RING_BUFFER_SIZE: usize = 192_000 * 2;
29
30const SEEK_END_GUARD_MS: u64 = 500;
33
34#[derive(Debug, Error)]
35pub enum PlayerError {
36 #[error("backend error: {0}")]
37 Backend(#[from] BackendError),
38 #[error("decode error: {0}")]
39 Decode(#[from] buffer::DecodeError),
40}
41
42#[derive(Clone)]
45struct StreamSource {
46 path: PathBuf,
47 bytes_written: Arc<crate::remote::downloads::ByteFeed>,
48 total: u64,
49 mode: streaming::ProbeMode,
50}
51
52fn media_extension(path: &Path) -> Option<String> {
55 crate::remote::download::strip_part_suffix(path)
56 .extension()
57 .and_then(|e| e.to_str())
58 .map(str::to_ascii_lowercase)
59}
60
61fn hint_for(path: &Path) -> symphonia::core::formats::probe::Hint {
63 let mut hint = symphonia::core::formats::probe::Hint::new();
64 if let Some(ext) = media_extension(path) {
65 hint.with_extension(&ext);
66 }
67 hint
68}
69
70fn lengthless_mode_for(path: &Path) -> streaming::ProbeMode {
80 match media_extension(path).as_deref() {
81 Some("ogg" | "oga" | "opus" | "spx" | "m4a" | "m4b" | "mp4" | "mov") => {
82 streaming::ProbeMode::LengthlessWholeEnd
83 }
84 _ => streaming::ProbeMode::Lengthless,
85 }
86}
87
88pub struct Player {
90 shared_state: Arc<SharedPlayerState>,
91 commands: CommandChannel,
92 active_playback: Option<ActivePlayback>,
93 timeline: Arc<PlaybackTimeline>,
94 viz_buffer: Arc<VizBuffer>,
95 viz_snapshot: Arc<VizSnapshot>,
96 _viz_analyzer: VizAnalyzer,
98 undo_stack: UndoStack,
99 batch_buffer: Option<Vec<UndoEntry>>,
102 output_device_name: Option<String>,
104 backend: Box<dyn AudioBackend>,
106 last_skip: std::time::Instant,
108 stream_mode: streaming::ProbeMode,
111 history: Option<PlayRecorder>,
114 in_flight: Option<InFlight>,
116 #[cfg(test)]
119 playback_starts: usize,
120}
121
122struct ActivePlayback {
124 engine: Box<dyn AudioEngineHandle>,
125 decode_handle: buffer::DecodeHandle,
126 stream: Option<LiveStream>,
128 _rate_watch: Option<Box<dyn SampleRateWatch>>,
131}
132
133struct LiveStream {
137 feed: Arc<crate::remote::downloads::ByteFeed>,
138 abandoned: Arc<std::sync::atomic::AtomicBool>,
139}
140
141impl LiveStream {
142 fn abandon(&self) {
143 self.abandoned
144 .store(true, std::sync::atomic::Ordering::Release);
145 self.feed.done();
146 }
147}
148
149impl Default for Player {
150 fn default() -> Self {
151 Self::new()
152 }
153}
154
155impl Player {
156 pub fn new() -> Self {
157 let viz_buffer = VizBuffer::new();
158 let viz_snapshot = VizSnapshot::new();
159 let timeline = PlaybackTimeline::new();
160 let cfg = crate::config::Config::load_or_default();
161 let viz_analyzer = VizAnalyzer::spawn_with_snapshot(
162 Arc::clone(&viz_buffer),
163 &cfg.visualizer,
164 Arc::clone(&viz_snapshot),
165 timeline.samples_played_counter(),
166 );
167
168 Self {
169 shared_state: SharedPlayerState::new(),
170 commands: CommandChannel::new(),
171 active_playback: None,
172 timeline,
173 viz_buffer,
174 viz_snapshot,
175 _viz_analyzer: viz_analyzer,
176 undo_stack: UndoStack::new(),
177 batch_buffer: None,
178 output_device_name: cfg.playback.output_device,
179 backend: crate::audio::platform_backend(),
180 last_skip: std::time::Instant::now(),
181 stream_mode: streaming::ProbeMode::Full,
182 history: None,
183 in_flight: None,
184 #[cfg(test)]
185 playback_starts: 0,
186 }
187 }
188
189 pub fn shared_state(&self) -> Arc<SharedPlayerState> {
191 self.shared_state.clone()
192 }
193
194 pub fn timeline(&self) -> Arc<PlaybackTimeline> {
196 self.timeline.clone()
197 }
198
199 pub fn viz_buffer(&self) -> Arc<VizBuffer> {
201 self.viz_buffer.clone()
202 }
203
204 pub fn viz_snapshot(&self) -> Arc<VizSnapshot> {
207 self.viz_snapshot.clone()
208 }
209
210 pub fn undo_stack(&self) -> &UndoStack {
212 &self.undo_stack
213 }
214
215 #[allow(clippy::type_complexity)]
223 fn create_engine_for(
224 &self,
225 info: &buffer::StreamInfo,
226 consumer: rtrb::Consumer<f32>,
227 ) -> Result<(Box<dyn AudioEngineHandle>, Option<Box<dyn SampleRateWatch>>), PlayerError> {
228 let device = self.resolve_device()?;
229 let device_rate = self.backend.get_device_sample_rate(&device)?;
230 let source_rate = info.sample_rate as f64;
231
232 self.shared_state.clear_output_sample_rate();
236
237 let settled = if (device_rate - source_rate).abs() > 0.1 {
238 log::info!(
239 "switching device sample rate: {}Hz → {}Hz",
240 device_rate,
241 source_rate
242 );
243 match self.backend.set_device_sample_rate(&device, source_rate) {
244 Ok(rate) => rate,
245 Err(e) => {
246 log::warn!("failed to set device sample rate: {}", e);
247 device_rate
248 }
249 }
250 } else {
251 device_rate
252 };
253
254 if (settled - source_rate).abs() > 0.1 {
255 log::warn!(
256 "device stayed at {}Hz (wanted {}Hz) — output is resampled, not bit-perfect",
257 settled,
258 source_rate
259 );
260 }
261
262 self.shared_state
265 .set_output_sample_rate(settled.round() as u32);
266
267 let watch_state = self.shared_state.clone();
271 let watch_name = device.name.clone();
272 let rate_watch = self.backend.watch_device_sample_rate(
273 &device,
274 Box::new(move |rate| {
275 log::info!("device sample rate changed externally: {rate}Hz on '{watch_name}'");
276 watch_state.set_output_sample_rate(rate.round() as u32);
277 }),
278 );
279
280 let engine = self.backend.create_engine(
281 &device,
282 source_rate,
283 info.channels as u32,
284 consumer,
285 self.timeline.samples_played_counter(),
286 )?;
287
288 Ok((engine, rate_watch))
289 }
290
291 fn resolve_device(&self) -> Result<backend::DeviceInfo, PlayerError> {
294 if let Some(ref name) = self.output_device_name {
295 match self.backend.list_devices() {
296 Ok(devices) => {
297 if let Some(dev) = devices.into_iter().find(|d| d.name == *name) {
298 return Ok(dev);
299 }
300 log::warn!(
301 "configured output device '{}' not found, falling back to default",
302 name,
303 );
304 }
305 Err(e) => {
306 log::warn!("failed to list devices while resolving '{}': {}", name, e);
307 }
308 }
309 }
310 Ok(self.backend.default_device()?)
311 }
312
313 pub fn set_output_device(&mut self, name: String) {
316 log::info!("switching output device to: {}", name);
317 self.output_device_name = Some(name.clone());
318
319 if let Err(e) = crate::config::Config::persist(|cfg| {
320 cfg.playback.output_device = Some(name);
321 }) {
322 log::error!("failed to save output device config: {}", e);
323 }
324
325 self.restart_on_current_track();
326 }
327
328 pub fn clear_output_device(&mut self) {
330 log::info!("reverting to system default output device");
331 self.output_device_name = None;
332
333 if let Err(e) = crate::config::Config::persist(|cfg| {
334 cfg.playback.output_device = None;
335 }) {
336 log::error!("failed to save output device config: {}", e);
337 }
338
339 self.restart_on_current_track();
340 }
341
342 fn restart_on_current_track(&mut self) {
345 if let Some(info) = self.shared_state.track_info() {
346 let position_ms = self.shared_state.position_ms();
347 if let Err(e) = self.restart_current(&info, position_ms) {
348 log::error!("failed to restart playback on device switch: {}", e);
349 }
350 }
351 }
352
353 pub fn output_device_name(&self) -> Option<&str> {
355 self.output_device_name.as_deref()
356 }
357
358 pub fn command_sender(&self) -> crossbeam_channel::Sender<PlayerCommand> {
360 self.commands.tx.clone()
361 }
362
363 pub fn play(&mut self, id: QueueItemId) {
366 self.shared_state.set_cursor(Some(id));
367
368 match self.shared_state.item_playback_source(id) {
369 Some(PlaybackSource::Ready(path)) => {
370 if let Err(e) = self.start_playback(id, &path, 0) {
371 log::error!("play failed: {}", e);
372 }
373 }
374 Some(PlaybackSource::Streaming {
375 path,
376 bytes_written,
377 total,
378 }) => {
379 self.report(PlaybackReportState::Stopped);
383 self.stop_engine();
384 self.shared_state.set_playback_state(PlaybackState::Stopped);
385 self.probe_stream_for_playback(id, &path, bytes_written, total);
386 }
387 None => {
388 self.report(PlaybackReportState::Stopped);
390 self.stop_engine();
391 self.shared_state.set_playback_state(PlaybackState::Stopped);
392 log::info!("play: item {:?} not ready, waiting for TrackReady", id);
393 }
394 }
395 }
396
397 fn start_playback(
402 &mut self,
403 id: QueueItemId,
404 path: &Path,
405 seek_ms: u64,
406 ) -> Result<(), PlayerError> {
407 #[cfg(test)]
408 {
409 self.playback_starts += 1;
410 }
411 let result = self.open_playback(id, path, seek_ms);
412 if result.is_err() {
413 self.stop_playback_and_clear_state();
414 }
415 self.wake_analyzer();
416 result
417 }
418
419 fn open_playback(
420 &mut self,
421 id: QueueItemId,
422 path: &Path,
423 seek_ms: u64,
424 ) -> Result<(), PlayerError> {
425 self.stop_engine();
426
427 let info = buffer::probe_file(path)?;
428
429 self.shared_state.set_track_info(Some(TrackInfo {
433 id,
434 path: path.to_path_buf(),
435 codec: info.codec.clone(),
436 sample_rate: info.sample_rate,
437 bit_depth: info.bit_depth,
438 bitrate_kbps: info.bitrate_kbps,
439 channels: info.channels,
440 duration_ms: info.duration_ms,
441 }));
442 self.shared_state.set_position_ms(seek_ms);
443 self.on_track_changed(id, seek_ms);
444 log::info!(
445 "playing: {} ({:?}) — {} {}Hz/{}ch, {}ms{}",
446 path.display(),
447 id,
448 info.codec,
449 info.sample_rate,
450 info.channels,
451 info.duration_ms,
452 if seek_ms > 0 {
453 format!(" @{}ms", seek_ms)
454 } else {
455 String::new()
456 }
457 );
458
459 let (producer, consumer) = rtrb::RingBuffer::new(RING_BUFFER_SIZE);
460
461 self.timeline.reset();
462
463 let next_track = self.decode_cursor(id);
464
465 let cfg = crate::config::Config::load_or_default();
467 let rg_mode = cfg.playback.replaygain;
468 let pre_amp_db = cfg.playback.pre_amp_db;
469
470 let finish_tx = self.commands.tx.clone();
471 let (_stream_info, decode_handle) = buffer::start_decode_file(
472 id,
473 path,
474 producer,
475 seek_ms,
476 next_track,
477 self.timeline.clone(),
478 Some(self.viz_buffer.clone()),
479 rg_mode,
480 pre_amp_db,
481 move || {
482 finish_tx.send(PlayerCommand::DecodeFinished).ok();
483 },
484 )?;
485
486 let (engine, rate_watch) = self.create_engine_for(&info, consumer)?;
487 engine.start()?;
488
489 self.shared_state.set_playback_state(PlaybackState::Playing);
490
491 self.active_playback = Some(ActivePlayback {
492 engine,
493 decode_handle,
494 stream: None,
495 _rate_watch: rate_watch,
496 });
497
498 Ok(())
499 }
500
501 fn decode_cursor(
505 &self,
506 id: QueueItemId,
507 ) -> impl Fn() -> Option<(QueueItemId, PathBuf)> + Send + 'static {
508 let state = self.shared_state.clone();
509 let cursor = parking_lot::Mutex::new(Some(id));
510 move || {
511 let current = cursor.lock().take()?;
512 let next = state.peek_next_ready_after(current);
513 if let Some((next_id, _)) = &next {
514 *cursor.lock() = Some(*next_id);
515 }
516 next
517 }
518 }
519
520 fn probe_stream_for_playback(
535 &self,
536 id: QueueItemId,
537 path: &Path,
538 bytes_written: Arc<crate::remote::downloads::ByteFeed>,
539 total: u64,
540 ) {
541 let path = path.to_path_buf();
542 let tx = self.commands.tx.clone();
543 let hint = hint_for(&path);
544
545 let status = {
550 let downloading = self.stream_status_fn(id);
551 let state = self.shared_state.clone();
552 Arc::new(move || {
553 if state.is_cursor(id) {
554 downloading()
555 } else {
556 streaming::StreamStatus::Failed
557 }
558 }) as Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync>
559 };
560
561 let spawned = thread::Builder::new()
562 .name("koan-stream-probe".into())
563 .spawn(move || {
564 let attempt = |mode, wait: bool| {
571 let open = if wait {
572 streaming::PartialFileSource::open(
573 &path,
574 bytes_written.clone(),
575 total,
576 status.clone(),
577 mode,
578 )
579 } else {
580 streaming::PartialFileSource::open_for_probe(
581 &path,
582 bytes_written.clone(),
583 total,
584 status.clone(),
585 mode,
586 )
587 };
588 open.map_err(buffer::DecodeError::Io).and_then(|source| {
589 let mss = symphonia::core::io::MediaSourceStream::new(
590 Box::new(source),
591 Default::default(),
592 );
593 buffer::probe_source(mss, &hint)
594 })
595 };
596
597 let info = match attempt(streaming::ProbeMode::Full, false) {
601 Ok(info) => Some((info, streaming::ProbeMode::Full)),
602 Err(e) => {
603 log::info!(
609 "stream probe: {} needs more than has arrived ({}), opening without a length",
610 path.display(),
611 e
612 );
613 let lengthless = lengthless_mode_for(&path);
614 attempt(lengthless, true)
615 .ok()
616 .map(|info| (info, lengthless))
617 }
618 };
619
620 match info {
621 Some((info, mode)) => {
622 tx.send(PlayerCommand::StreamProbed {
623 id,
624 info: Box::new(info),
625 mode,
626 })
627 .ok();
628 }
629 None => log::info!(
632 "stream probe: {} cannot start early, waiting for the download",
633 path.display()
634 ),
635 }
636 });
637
638 if let Err(e) = spawned {
639 log::warn!("stream probe: could not spawn for {:?}: {}", id, e);
640 }
641 }
642
643 fn stream_probed(
646 &mut self,
647 id: QueueItemId,
648 info: buffer::StreamInfo,
649 mode: streaming::ProbeMode,
650 ) {
651 if !self.shared_state.is_cursor(id) {
652 return; }
654 if self.shared_state.playback_state() != PlaybackState::Stopped {
655 return; }
657
658 match self.shared_state.item_playback_source(id) {
659 Some(PlaybackSource::Ready(path)) => {
661 if let Err(e) = self.start_playback(id, &path, 0) {
662 log::error!("stream probe: playback failed: {}", e);
663 }
664 }
665 Some(PlaybackSource::Streaming {
666 path,
667 bytes_written,
668 total,
669 }) => {
670 let source = StreamSource {
671 path,
672 bytes_written,
673 total,
674 mode,
675 };
676 if let Err(e) = self.start_streaming_playback(id, source, 0, info) {
677 log::error!("stream probe: streaming playback failed: {}", e);
678 }
679 }
680 None => {}
681 }
682 }
683
684 fn stream_status_fn(
688 &self,
689 id: QueueItemId,
690 ) -> Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync> {
691 let state = self.shared_state.clone();
692 Arc::new(move || match state.item_load_state(id) {
693 Some(LoadState::Ready) => streaming::StreamStatus::Complete,
694 Some(LoadState::Failed(_)) => streaming::StreamStatus::Failed,
695 _ => streaming::StreamStatus::Downloading,
696 })
697 }
698
699 fn start_streaming_playback(
709 &mut self,
710 id: QueueItemId,
711 source: StreamSource,
712 seek_ms: u64,
713 info: buffer::StreamInfo,
714 ) -> Result<(), PlayerError> {
715 let result = self.open_streaming_playback(id, source, seek_ms, info);
716 if result.is_err() {
717 self.stop_playback_and_clear_state();
718 }
719 result
720 }
721
722 fn open_streaming_playback(
723 &mut self,
724 id: QueueItemId,
725 source: StreamSource,
726 seek_ms: u64,
727 info: buffer::StreamInfo,
728 ) -> Result<(), PlayerError> {
729 self.stop_engine();
730 self.stream_mode = source.mode;
732 let path = source.path.as_path();
733
734 let live = LiveStream {
735 feed: source.bytes_written.clone(),
736 abandoned: Default::default(),
737 };
738 let status = {
739 let downloading = self.stream_status_fn(id);
740 let abandoned = live.abandoned.clone();
741 Arc::new(move || {
742 if abandoned.load(std::sync::atomic::Ordering::Acquire) {
743 streaming::StreamStatus::Failed
744 } else {
745 downloading()
746 }
747 }) as Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync>
748 };
749 let open_source = {
750 let StreamSource {
751 path,
752 bytes_written,
753 total,
754 mode,
755 } = source.clone();
756 let status = status.clone();
757 move || {
758 streaming::PartialFileSource::open(
759 &path,
760 bytes_written.clone(),
761 total,
762 status.clone(),
763 mode,
764 )
765 }
766 };
767
768 self.shared_state.set_track_info(Some(TrackInfo {
769 id,
770 path: path.to_path_buf(),
771 codec: info.codec.clone(),
772 sample_rate: info.sample_rate,
773 bit_depth: info.bit_depth,
774 bitrate_kbps: info.bitrate_kbps,
775 channels: info.channels,
776 duration_ms: info.duration_ms,
777 }));
778 self.shared_state.set_position_ms(seek_ms);
779 self.on_track_changed(id, seek_ms);
780 log::info!(
781 "streaming: {} ({:?}) — {} {}Hz/{}ch, {}ms{}",
782 path.display(),
783 id,
784 info.codec,
785 info.sample_rate,
786 info.channels,
787 info.duration_ms,
788 if seek_ms > 0 {
789 format!(" @{}ms", seek_ms)
790 } else {
791 String::new()
792 },
793 );
794
795 let (producer, consumer) = rtrb::RingBuffer::new(RING_BUFFER_SIZE);
796
797 self.timeline.reset();
798
799 let next_track = self.decode_cursor(id);
801
802 let first = buffer::SourceEntry {
803 id,
804 path: path.to_path_buf(),
805 hint: hint_for(path),
806 make_mss: Box::new(move || {
807 Ok(symphonia::core::io::MediaSourceStream::new(
808 Box::new(open_source()?),
809 Default::default(),
810 ))
811 }),
812 };
813
814 let cfg = crate::config::Config::load_or_default();
816 let rg_mode = cfg.playback.replaygain;
817 let pre_amp_db = cfg.playback.pre_amp_db;
818
819 let finish_tx = self.commands.tx.clone();
820 let (_stream_info, decode_handle) = buffer::start_decode(
821 first,
822 producer,
823 seek_ms,
824 move || {
825 let (next_id, next_path) = next_track()?;
826 Some(buffer::SourceEntry::from_file(next_id, next_path))
827 },
828 self.timeline.clone(),
829 Some(self.viz_buffer.clone()),
830 rg_mode,
831 pre_amp_db,
832 move || {
833 finish_tx.send(PlayerCommand::DecodeFinished).ok();
834 },
835 )?;
836
837 let (engine, rate_watch) = self.create_engine_for(&info, consumer)?;
838 engine.start()?;
839
840 self.shared_state.set_playback_state(PlaybackState::Playing);
841
842 self.active_playback = Some(ActivePlayback {
843 engine,
844 decode_handle,
845 stream: Some(live),
846 _rate_watch: rate_watch,
847 });
848
849 Ok(())
850 }
851
852 pub fn seek(&mut self, position_ms: u64) {
859 let Some(info) = self.shared_state.track_info() else {
860 return;
861 };
862 let seekable = self.shared_state.seekable_ms();
864 if seekable == 0 {
865 log::debug!("seek declined: {:?} is not seekable yet", info.id);
869 return;
870 }
871 let ceiling = seekable.min(
872 self.shared_state
873 .duration_ms()
874 .saturating_sub(SEEK_END_GUARD_MS),
875 );
876 let clamped = position_ms.min(ceiling);
877
878 if let Err(e) = self.restart_current(&info, clamped) {
879 log::error!("seek failed: {}", e);
880 }
881 }
882
883 fn restart_current(&mut self, info: &TrackInfo, position_ms: u64) -> Result<(), PlayerError> {
891 let was_paused = self.shared_state.playback_state() == PlaybackState::Paused;
892
893 match self.shared_state.item_playback_source(info.id) {
894 Some(PlaybackSource::Streaming {
895 path,
896 bytes_written,
897 total,
898 }) => {
899 let known = buffer::StreamInfo {
903 codec: info.codec.clone(),
904 sample_rate: info.sample_rate,
905 channels: info.channels,
906 bit_depth: info.bit_depth,
907 bitrate_kbps: info.bitrate_kbps,
908 duration_ms: info.duration_ms,
909 };
910 let source = StreamSource {
911 path,
912 bytes_written,
913 total,
914 mode: self.stream_mode,
915 };
916 self.start_streaming_playback(info.id, source, position_ms, known)?;
917 }
918 Some(PlaybackSource::Ready(path)) => {
919 self.start_playback(info.id, &path, position_ms)?;
920 }
921 None => return Ok(()),
922 }
923
924 if was_paused {
925 self.pause_now();
926 }
927 self.report(if was_paused {
928 PlaybackReportState::Paused
929 } else {
930 PlaybackReportState::Playing
931 });
932 Ok(())
933 }
934
935 pub fn next_track(&mut self) {
937 match self.shared_state.advance_cursor_loadable() {
938 Some(id) => self.play(id),
939 None => {
940 log::info!("no more tracks in playlist");
941 self.stop_playback_and_clear_state();
942 }
943 }
944 }
945
946 pub fn prev_track(&mut self) {
948 match self.shared_state.retreat_cursor() {
949 Some((id, _)) => self.play(id),
950 None => {
951 if let Some(info) = self.shared_state.track_info()
953 && let Err(e) = self.restart_current(&info, 0)
954 {
955 log::error!("restart failed: {}", e);
956 }
957 }
958 }
959 }
960
961 pub fn pause(&mut self) {
966 let Some(ref playback) = self.active_playback else {
967 return;
968 };
969 if crate::config::Config::load_or_default()
970 .playback
971 .fade_on_pause
972 {
973 playback.engine.fade_out();
974 self.shared_state.set_playback_state(PlaybackState::Paused);
975 } else {
976 self.pause_now();
977 }
978 self.report(PlaybackReportState::Paused);
979 }
980
981 fn pause_now(&mut self) {
984 if let Some(ref playback) = self.active_playback {
985 if let Err(e) = playback.engine.stop() {
986 log::error!("pause failed: {}", e);
987 return;
988 }
989 self.shared_state.set_playback_state(PlaybackState::Paused);
990 }
991 }
992
993 pub fn resume(&mut self) {
999 if self.active_playback.is_none() {
1000 if let Some(id) = self.shared_state.cursor() {
1001 self.play(id);
1002 }
1003 return;
1004 }
1005 if let Some(ref playback) = self.active_playback {
1006 let engine = &playback.engine;
1007 let resumed = if engine.is_running() || engine.is_silent() {
1008 engine.fade_in()
1009 } else {
1010 engine.start()
1011 };
1012 if let Err(e) = resumed {
1013 log::error!("resume failed: {}", e);
1014 return;
1015 }
1016 self.shared_state.set_playback_state(PlaybackState::Playing);
1017 self.wake_analyzer();
1018 self.report(PlaybackReportState::Playing);
1019 }
1020 }
1021
1022 fn wake_analyzer(&self) {
1029 self.viz_snapshot.wake();
1030 }
1031
1032 pub fn stop(&mut self) {
1034 self.shared_state.clear_playlist();
1035 self.stop_playback_and_clear_state();
1036 }
1037
1038 fn stop_engine(&mut self) {
1044 let Some(playback) = self.active_playback.take() else {
1045 return;
1046 };
1047 let ActivePlayback {
1048 engine,
1049 mut decode_handle,
1050 stream,
1051 _rate_watch,
1052 } = playback;
1053
1054 let _ = engine.stop();
1055 decode_handle.signal_stop();
1058 if let Some(stream) = stream {
1059 stream.abandon();
1060 }
1061 decode_handle.stop();
1062 drop(engine);
1063 }
1064
1065 fn stop_playback_and_clear_state(&mut self) {
1067 self.report(PlaybackReportState::Stopped);
1068 self.finish_play();
1069 self.stop_engine();
1070 self.timeline.reset();
1071 self.shared_state.set_playback_state(PlaybackState::Stopped);
1072 self.shared_state.set_position_ms(0);
1073 self.shared_state.set_track_info(None);
1074 }
1075
1076 pub fn remove_from_playlist(&mut self, id: QueueItemId) {
1083 let was_cursor = self.shared_state.is_cursor(id);
1084 let resume_after = was_cursor
1085 .then(|| self.shared_state.item_before(id))
1086 .flatten();
1087 self.shared_state.remove_item(id);
1088 if was_cursor {
1089 self.shared_state.set_cursor(resume_after);
1090 self.next_track();
1091 }
1092 }
1093
1094 pub fn track_ready(&mut self, id: QueueItemId) {
1097 self.shared_state.update_item_state(id, ItemState::Ready);
1099
1100 if !self.shared_state.is_cursor(id) {
1101 return;
1102 }
1103
1104 let is_playing = self.shared_state.playback_state() == PlaybackState::Playing;
1105 let current_track_id = self.shared_state.track_info().map(|t| t.id);
1106
1107 if is_playing && current_track_id == Some(id) {
1108 log::info!(
1111 "track_ready: download complete while streaming {:?}, refreshing metadata",
1112 id
1113 );
1114 self.refresh_track_metadata(id);
1115 return;
1116 }
1117
1118 if !is_playing && let Some(path) = self.shared_state.item_path_if_ready(id) {
1120 log::info!("track_ready: starting playback for {:?}", id);
1121 if let Err(e) = self.start_playback(id, &path, 0) {
1122 log::error!("track_ready playback failed: {}", e);
1123 }
1124 }
1125 }
1126
1127 pub fn track_stream_ready(&mut self, id: QueueItemId) {
1130 if !self.shared_state.is_cursor(id) {
1131 return;
1132 }
1133
1134 let is_playing = self.shared_state.playback_state() == PlaybackState::Playing;
1135 if is_playing {
1136 return; }
1138
1139 match self.shared_state.item_playback_source(id) {
1140 Some(PlaybackSource::Streaming {
1141 path,
1142 bytes_written,
1143 total,
1144 }) => {
1145 log::info!("track_stream_ready: probing partial file for {:?}", id);
1146 self.probe_stream_for_playback(id, &path, bytes_written, total);
1147 }
1148 Some(PlaybackSource::Ready(path)) => {
1149 log::info!(
1151 "track_stream_ready: track already ready, starting normal playback for {:?}",
1152 id
1153 );
1154 if let Err(e) = self.start_playback(id, &path, 0) {
1155 log::error!("track_stream_ready playback failed: {}", e);
1156 }
1157 }
1158 None => {} }
1160 }
1161
1162 fn refresh_track_metadata(&mut self, id: QueueItemId) {
1166 use crate::index::metadata;
1167
1168 let path = match self.shared_state.item_path_if_ready(id) {
1169 Some(p) => p,
1170 None => return,
1171 };
1172
1173 match metadata::read_metadata(&path) {
1174 Ok(meta) => {
1175 self.shared_state.update_item_metadata(
1176 id,
1177 meta.title,
1178 meta.artist,
1179 meta.album_artist.unwrap_or_default(),
1180 meta.album,
1181 meta.duration_ms.map(|d| d as u64),
1182 );
1183
1184 if let Some(current) = self.shared_state.track_info()
1193 && current.id == id
1194 {
1195 let probed = buffer::probe_file(&path).ok();
1196 let duration_ms = probed
1197 .as_ref()
1198 .map(|s| s.duration_ms)
1199 .filter(|d| *d > current.duration_ms)
1200 .unwrap_or(current.duration_ms);
1201 if duration_ms != current.duration_ms {
1202 log::info!(
1203 "track_ready: duration corrected {}ms → {}ms",
1204 current.duration_ms,
1205 duration_ms
1206 );
1207 }
1208 self.shared_state.set_track_info(Some(TrackInfo {
1209 duration_ms,
1210 path: path.clone(),
1211 ..current
1212 }));
1213 }
1214
1215 self.shared_state.signal_metadata_refresh();
1217 log::info!("track_ready: metadata refreshed for {:?}", id);
1218 }
1219 Err(e) => {
1220 log::warn!("track_ready: metadata refresh failed for {:?}: {}", id, e);
1221 }
1222 }
1223 }
1224
1225 fn on_track_changed(&mut self, id: QueueItemId, position_ms: u64) {
1233 if self.in_flight.as_ref().is_some_and(|f| f.item == id) {
1234 return;
1235 }
1236 self.finish_play();
1237 let track_id = self.shared_state.item_db_id(id);
1238 self.in_flight = Some(InFlight::new(id, track_id));
1239 if let (Some(track_id), Some(recorder)) = (track_id, self.history.as_ref()) {
1240 recorder.record(PlayEvent::Started {
1241 track_id,
1242 position_ms,
1243 });
1244 }
1245 }
1246
1247 fn report(&self, state: PlaybackReportState) {
1251 let Some(track_id) = self.in_flight.as_ref().and_then(InFlight::track_id) else {
1252 return;
1253 };
1254 if let Some(recorder) = self.history.as_ref() {
1255 recorder.record(PlayEvent::Playback(PlaybackReport {
1256 track_id,
1257 state,
1258 position_ms: self.shared_state.position_ms(),
1259 }));
1260 }
1261 }
1262
1263 fn finish_play(&mut self) -> Option<PlayEvent> {
1266 let flight = self.in_flight.take()?;
1267 let event = PlayEvent::Finished {
1268 track_id: flight.track_id()?,
1269 listened_ms: flight.listened_ms(),
1270 };
1271 if let Some(recorder) = self.history.as_ref() {
1272 recorder.record(event);
1273 }
1274 Some(event)
1275 }
1276
1277 pub fn update_playback_state(&mut self) {
1280 let Some(playback) = self.active_playback.as_ref() else {
1281 return;
1282 };
1283
1284 if self.shared_state.playback_state() == PlaybackState::Paused
1285 && playback.engine.is_running()
1286 && playback.engine.is_silent()
1287 && let Err(e) = playback.engine.stop()
1288 {
1289 log::error!("stopping after fade failed: {}", e);
1290 }
1291
1292 if let Some((id, path, info, position_ms)) = self.timeline.current_playback() {
1293 self.shared_state.set_position_ms(position_ms);
1294
1295 self.on_track_changed(id, position_ms);
1298 if let Some(f) = self.in_flight.as_mut() {
1299 f.advance(position_ms);
1300 }
1301
1302 let current_id = self.shared_state.track_info().map(|t| t.id);
1305 if current_id != Some(id) {
1306 log::info!("timeline: now playing {:?}", id);
1307 self.shared_state.set_track_info(Some(TrackInfo {
1308 id,
1309 path,
1310 codec: info.codec,
1311 sample_rate: info.sample_rate,
1312 bit_depth: info.bit_depth,
1313 bitrate_kbps: info.bitrate_kbps,
1314 channels: info.channels,
1315 duration_ms: info.duration_ms,
1316 }));
1317 self.shared_state.set_cursor(Some(id));
1318 }
1319 }
1320 }
1321
1322 pub fn track_failed(&mut self, id: QueueItemId) {
1329 if !self.shared_state.is_cursor(id) {
1330 return;
1331 }
1332 if self.shared_state.playback_state() != PlaybackState::Stopped {
1336 return;
1337 }
1338 log::info!("track {:?} cannot load, moving on", id);
1339 self.next_track();
1340 }
1341
1342 fn on_decode_finished(&mut self) {
1349 log::info!("decode finished, checking for next track");
1350 match self.shared_state.advance_cursor_loadable() {
1351 Some(id) => self.play(id),
1352 None => {
1353 log::info!("no more tracks — stopping");
1354 self.stop_playback_and_clear_state();
1355 }
1356 }
1357 }
1358
1359 fn snapshot_for_undo(
1363 &self,
1364 ids: &[QueueItemId],
1365 ) -> Vec<(Box<state::PlaylistItem>, Option<QueueItemId>)> {
1366 self.shared_state
1367 .items_before(ids)
1368 .into_iter()
1369 .filter_map(|(id, after)| Some((Box::new(self.shared_state.get_item(id)?), after)))
1370 .collect()
1371 }
1372
1373 fn push_undo(&mut self, entry: UndoEntry) {
1375 if let Some(ref mut batch) = self.batch_buffer {
1376 batch.push(entry);
1377 } else {
1378 self.undo_stack.push(entry);
1379 }
1380 }
1381
1382 pub fn process_command(&mut self, cmd: PlayerCommand) {
1384 match cmd {
1385 PlayerCommand::Play(id) => self.play(id),
1386 PlayerCommand::Pause => self.pause(),
1387 PlayerCommand::Resume => self.resume(),
1388 PlayerCommand::Stop => self.stop(),
1389 PlayerCommand::Seek(pos) => self.seek(pos),
1390 PlayerCommand::NextTrack => {
1391 let now = std::time::Instant::now();
1393 if now.duration_since(self.last_skip).as_millis() >= 150 {
1394 self.last_skip = now;
1395 self.next_track();
1396 }
1397 }
1398 PlayerCommand::PrevTrack => {
1399 let now = std::time::Instant::now();
1400 if now.duration_since(self.last_skip).as_millis() >= 150 {
1401 self.last_skip = now;
1402 self.prev_track();
1403 }
1404 }
1405 PlayerCommand::AddToPlaylist(items) => {
1406 let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
1407 self.shared_state.add_items(items);
1408 self.push_undo(UndoEntry::Added { ids });
1409 }
1410 PlayerCommand::UpdatePaths(updates) => {
1411 self.shared_state.update_paths(&updates);
1412 if let Some(info) = self.shared_state.track_info()
1413 && let Some((_, new_path)) = updates.iter().find(|(id, _)| *id == info.id)
1414 {
1415 self.shared_state.set_track_info(Some(TrackInfo {
1416 path: new_path.clone(),
1417 ..info
1418 }));
1419 }
1420 }
1421 PlayerCommand::InsertInPlaylist { items, after } => {
1422 let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
1423 self.shared_state.insert_items_after(items, after);
1424 self.push_undo(UndoEntry::Inserted { ids });
1425 }
1426 PlayerCommand::ClearPlaylist => {
1427 self.stop_playback_and_clear_state();
1431 let (items, cursor) = self.shared_state.snapshot_playlist();
1432 self.shared_state.clear_playlist();
1433 self.push_undo(UndoEntry::Replaced { items, cursor });
1434 }
1435 PlayerCommand::ReplacePlaylist { items, start } => {
1436 self.stop_playback_and_clear_state();
1440 let (old_items, cursor) = self.shared_state.snapshot_playlist();
1441 self.shared_state.clear_playlist();
1442 self.push_undo(UndoEntry::Replaced {
1443 items: old_items,
1444 cursor,
1445 });
1446
1447 if items.is_empty() {
1448 return;
1449 }
1450 let start_id = items.get(start).unwrap_or(&items[0]).id;
1451 self.shared_state.add_items(items);
1452 self.play(start_id);
1453 }
1454 PlayerCommand::RemoveFromPlaylist(id) => {
1455 let item = self.shared_state.get_item(id);
1456 let after = self.shared_state.item_before(id);
1457 self.remove_from_playlist(id);
1458 if let Some(item) = item {
1459 self.push_undo(UndoEntry::Removed {
1460 items: vec![(Box::new(item), after)],
1461 });
1462 }
1463 }
1464 PlayerCommand::RemoveFromPlaylistBatch(ids) => {
1465 let items_with_pos = self.snapshot_for_undo(&ids);
1469 let resume_after = match self.shared_state.cursor() {
1470 Some(cursor) if ids.contains(&cursor) => {
1471 Some(self.shared_state.surviving_item_before(cursor, &ids))
1472 }
1473 _ => None,
1474 };
1475
1476 self.shared_state.remove_items(&ids);
1477
1478 if let Some(resume_after) = resume_after {
1479 self.shared_state.set_cursor(resume_after);
1480 self.next_track();
1481 }
1482
1483 if !items_with_pos.is_empty() {
1484 self.push_undo(UndoEntry::Removed {
1485 items: items_with_pos,
1486 });
1487 }
1488 }
1489 PlayerCommand::MoveInPlaylist { id, target, after } => {
1490 let was_after = self.shared_state.item_before(id);
1491 self.shared_state.move_item(id, target, after);
1492 self.push_undo(UndoEntry::Moved { id, was_after });
1493 }
1494 PlayerCommand::MoveItemsInPlaylist { ids, target, after } => {
1495 let entries = self.shared_state.items_before(&ids);
1496 self.shared_state.move_items(&ids, target, after);
1497 self.push_undo(UndoEntry::MovedBatch { entries });
1498 }
1499 PlayerCommand::ReorderPlaylist(order) => {
1500 let entries = self.shared_state.items_before(&order);
1504 self.shared_state.reorder_to(&order);
1505 self.push_undo(UndoEntry::MovedBatch { entries });
1506 }
1507 PlayerCommand::TrackReady(id) => self.track_ready(id),
1508 PlayerCommand::DecodeFinished => self.on_decode_finished(),
1509 PlayerCommand::TrackStreamReady(id) => self.track_stream_ready(id),
1510 PlayerCommand::StreamProbed { id, info, mode } => self.stream_probed(id, *info, mode),
1511 PlayerCommand::TrackFailed(id) => self.track_failed(id),
1512 PlayerCommand::Undo => self.execute_undo(),
1513 PlayerCommand::Redo => self.execute_redo(),
1514 PlayerCommand::BeginUndoBatch => {
1515 self.batch_buffer = Some(Vec::new());
1516 }
1517 PlayerCommand::EndUndoBatch => {
1518 if let Some(entries) = self.batch_buffer.take() {
1519 if entries.len() == 1 {
1520 self.undo_stack.push(entries.into_iter().next().unwrap());
1522 } else if !entries.is_empty() {
1523 self.undo_stack.push(UndoEntry::Batch(entries));
1524 }
1525 }
1526 }
1527 PlayerCommand::SetOutputDevice(name) => self.set_output_device(name),
1528 PlayerCommand::RestartOutput => {
1529 log::info!("restarting audio output");
1530 self.restart_on_current_track();
1531 }
1532 PlayerCommand::ClearOutputDevice => self.clear_output_device(),
1533 }
1534 }
1535
1536 fn apply_entry(&mut self, entry: UndoEntry) -> UndoEntry {
1538 match entry {
1539 UndoEntry::Added { ids } | UndoEntry::Inserted { ids } => {
1540 let items_with_pos = self.snapshot_for_undo(&ids);
1542 self.shared_state.remove_items(&ids);
1543 UndoEntry::Removed {
1544 items: items_with_pos,
1545 }
1546 }
1547 UndoEntry::Removed { items } => {
1548 let mut ids = Vec::with_capacity(items.len());
1550 for (item, after) in items {
1551 ids.push(item.id);
1552 self.shared_state.insert_item_at(*item, after);
1553 }
1554 UndoEntry::Added { ids }
1555 }
1556 UndoEntry::Moved { id, was_after } => {
1557 let current_after = self.shared_state.item_before(id);
1558 self.shared_state.move_item_to(id, was_after);
1559 UndoEntry::Moved {
1560 id,
1561 was_after: current_after,
1562 }
1563 }
1564 UndoEntry::MovedBatch { entries } => {
1565 let ids: Vec<QueueItemId> = entries.iter().map(|(id, _)| *id).collect();
1566 let current_positions = self.shared_state.items_before(&ids);
1567 self.shared_state.move_items_to(&entries);
1568 UndoEntry::MovedBatch {
1569 entries: current_positions,
1570 }
1571 }
1572 UndoEntry::Replaced { items, cursor } => {
1573 let (current_items, current_cursor) = self.shared_state.snapshot_playlist();
1574 self.shared_state.restore_playlist(items, cursor);
1575 UndoEntry::Replaced {
1576 items: current_items,
1577 cursor: current_cursor,
1578 }
1579 }
1580 UndoEntry::Batch(entries) => {
1581 let mut inverses: Vec<_> = entries
1583 .into_iter()
1584 .rev()
1585 .map(|e| self.apply_entry(e))
1586 .collect();
1587 inverses.reverse();
1588 UndoEntry::Batch(inverses)
1589 }
1590 }
1591 }
1592
1593 fn reconcile_playback(&mut self) {
1607 let Some(playing) = self.shared_state.track_info().map(|t| t.id) else {
1608 return;
1609 };
1610 if self.shared_state.get_item(playing).is_some() {
1611 return;
1612 }
1613 let resume = (self.shared_state.playback_state() == PlaybackState::Playing)
1618 .then(|| self.shared_state.cursor())
1619 .flatten();
1620 self.stop_playback_and_clear_state();
1621 if let Some(id) = resume {
1622 self.play(id);
1623 }
1624 }
1625
1626 fn execute_undo(&mut self) {
1628 let Some(entry) = self.undo_stack.pop_undo() else {
1629 return;
1630 };
1631 let inverse = self.apply_entry(entry);
1632 self.undo_stack.push_redo(inverse);
1633 self.reconcile_playback();
1634 }
1635
1636 fn execute_redo(&mut self) {
1638 let Some(entry) = self.undo_stack.pop_redo() else {
1639 return;
1640 };
1641 let inverse = self.apply_entry(entry);
1642 self.undo_stack.push_undo_keep_redo(inverse);
1643 self.reconcile_playback();
1644 }
1645
1646 pub fn run(&mut self) {
1648 use std::time::Duration;
1649
1650 let rx = self.commands.rx.clone();
1651 loop {
1652 match rx.recv_timeout(Duration::from_millis(50)) {
1654 Ok(cmd) => self.process_command(cmd),
1655 Err(crossbeam_channel::RecvTimeoutError::Timeout) => {}
1656 Err(crossbeam_channel::RecvTimeoutError::Disconnected) => break,
1657 }
1658 self.update_playback_state();
1659 }
1660 self.stop();
1661 }
1662
1663 pub fn spawn() -> (
1666 Arc<SharedPlayerState>,
1667 Arc<PlaybackTimeline>,
1668 Arc<VizSnapshot>,
1669 crossbeam_channel::Sender<PlayerCommand>,
1670 ) {
1671 let mut player = Self::new();
1672 player.history = PlayRecorder::spawn();
1673 let state = player.shared_state();
1674 let timeline = player.timeline();
1675 let viz_snapshot = player.viz_snapshot();
1676 let tx = player.command_sender();
1677
1678 thread::Builder::new()
1679 .name("koan-player".into())
1680 .spawn(move || player.run())
1681 .expect("failed to spawn player thread");
1682
1683 (state, timeline, viz_snapshot, tx)
1684 }
1685}
1686
1687#[cfg(test)]
1688mod tests {
1689 #[test]
1690 fn a_download_in_progress_is_known_by_its_own_extension() {
1691 use std::path::Path;
1692 assert_eq!(
1694 lengthless_mode_for(Path::new("/c/t.m4a.part")),
1695 streaming::ProbeMode::LengthlessWholeEnd
1696 );
1697 assert_eq!(
1698 lengthless_mode_for(Path::new("/c/t.OPUS.part")),
1699 streaming::ProbeMode::LengthlessWholeEnd
1700 );
1701 assert_eq!(
1702 lengthless_mode_for(Path::new("/c/t.flac.part")),
1703 streaming::ProbeMode::Lengthless
1704 );
1705 assert_eq!(
1706 media_extension(Path::new("/c/t.m4a.part")).as_deref(),
1707 Some("m4a")
1708 );
1709 assert_eq!(
1710 media_extension(Path::new("/c/t.mp3")).as_deref(),
1711 Some("mp3")
1712 );
1713 }
1714
1715 use super::*;
1716 use state::PlaylistItem;
1717 use std::path::PathBuf;
1718 use std::sync::atomic::AtomicU64;
1719
1720 fn make_item(title: &str) -> PlaylistItem {
1721 PlaylistItem {
1722 playlist_entry_id: None,
1723 id: QueueItemId::new(),
1724 db_id: None,
1725 path: PathBuf::from(format!("/music/{title}.flac")),
1726 title: title.to_string(),
1727 artist: String::new(),
1728 album_artist: String::new(),
1729 album: String::new(),
1730 year: None,
1731 codec: None,
1732 track_number: None,
1733 disc: None,
1734 duration_ms: None,
1735 state: ItemState::Ready,
1736 }
1737 }
1738
1739 fn playlist_ids(player: &Player) -> Vec<QueueItemId> {
1740 let (items, _) = player.shared_state.snapshot_playlist();
1741 items.iter().map(|i| i.id).collect()
1742 }
1743
1744 fn playlist_titles(player: &Player) -> Vec<String> {
1745 let (items, _) = player.shared_state.snapshot_playlist();
1746 items.iter().map(|i| i.title.clone()).collect()
1747 }
1748
1749 fn pending_item(title: &str) -> PlaylistItem {
1750 PlaylistItem {
1751 playlist_entry_id: None,
1752 state: ItemState::Pending,
1753 ..make_item(title)
1754 }
1755 }
1756
1757 fn pretend_playing(player: &mut Player, id: QueueItemId) {
1761 let item = player
1762 .shared_state
1763 .get_item(id)
1764 .expect("item is in the queue");
1765 player.shared_state.set_track_info(Some(TrackInfo {
1766 id,
1767 path: item.path,
1768 codec: String::new(),
1769 sample_rate: 44_100,
1770 bit_depth: None,
1771 bitrate_kbps: None,
1772 channels: 2,
1773 duration_ms: 1_000,
1774 }));
1775 player
1776 .shared_state
1777 .set_playback_state(PlaybackState::Playing);
1778 }
1779
1780 fn playing_id(player: &Player) -> Option<QueueItemId> {
1781 player.shared_state.track_info().map(|t| t.id)
1782 }
1783
1784 fn seed(player: &mut Player, n: usize) -> Vec<QueueItemId> {
1786 let items: Vec<_> = (0..n).map(|i| make_item(&format!("t{i}"))).collect();
1787 let ids = items.iter().map(|i| i.id).collect();
1788 player.process_command(PlayerCommand::AddToPlaylist(items));
1789 ids
1790 }
1791
1792 fn listen(player: &mut Player, from_ms: u64, to_ms: u64) {
1796 let mut at = from_ms;
1797 if let Some(f) = player.in_flight.as_mut() {
1798 f.advance(at); }
1800 while at < to_ms {
1801 at = (at + 50).min(to_ms);
1802 if let Some(f) = player.in_flight.as_mut() {
1803 f.advance(at);
1804 }
1805 }
1806 }
1807
1808 fn start(player: &mut Player, track_id: i64) -> QueueItemId {
1809 let id = QueueItemId::new();
1810 player.on_track_changed(id, 0);
1811 player
1813 .in_flight
1814 .as_mut()
1815 .unwrap()
1816 .track_id_for_test(track_id);
1817 id
1818 }
1819
1820 #[test]
1821 fn a_gapless_transition_closes_the_outgoing_track_and_opens_the_next() {
1822 let mut player = Player::new();
1823 start(&mut player, 11);
1824 listen(&mut player, 0, 200_000);
1825
1826 let b = QueueItemId::new();
1827 player.on_track_changed(b, 0);
1828 let f = player
1829 .in_flight
1830 .as_ref()
1831 .expect("the next track is counting");
1832 assert_eq!(f.item, b);
1833 assert_eq!(f.listened_ms(), 0, "and starts from nothing");
1834 }
1835
1836 #[test]
1837 fn a_track_skipped_seconds_in_is_still_history() {
1838 let mut player = Player::new();
1839 start(&mut player, 7);
1840 listen(&mut player, 0, 2_000);
1841
1842 let event = player
1843 .finish_play()
1844 .expect("putting something on is a thing you did, however briefly");
1845 assert!(matches!(
1846 event,
1847 history::PlayEvent::Finished {
1848 track_id: 7,
1849 listened_ms: 2_000
1850 }
1851 ));
1852 }
1853
1854 #[test]
1855 fn a_track_is_closed_out_once() {
1856 let mut player = Player::new();
1857 start(&mut player, 7);
1858 listen(&mut player, 0, 200_000);
1859
1860 assert!(player.finish_play().is_some());
1861 assert!(player.finish_play().is_none());
1862 }
1863
1864 #[test]
1865 fn seeking_around_a_track_does_not_enter_it_twice() {
1866 let mut player = Player::new();
1867 let id = start(&mut player, 7);
1868 listen(&mut player, 0, 120_000);
1869
1870 player.on_track_changed(id, 30_000);
1872 assert_eq!(
1873 player.in_flight.as_ref().unwrap().listened_ms(),
1874 120_000,
1875 "the seek kept the count rather than restarting it"
1876 );
1877 listen(&mut player, 30_000, 40_000);
1878
1879 let Some(history::PlayEvent::Finished { listened_ms, .. }) = player.finish_play() else {
1880 panic!("still one play");
1881 };
1882 assert_eq!(listened_ms, 130_000);
1883 assert!(player.finish_play().is_none());
1884 }
1885
1886 #[test]
1887 fn a_track_that_is_not_in_the_library_is_not_recorded() {
1888 let mut player = Player::new();
1889 let id = QueueItemId::new();
1890 player.on_track_changed(id, 0);
1891 listen(&mut player, 0, 200_000);
1892 assert!(player.finish_play().is_none());
1893 }
1894
1895 #[test]
1896 fn stopping_closes_out_what_was_heard() {
1897 let mut player = Player::new();
1898 start(&mut player, 7);
1899 listen(&mut player, 0, 150_000);
1900
1901 player.stop_playback_and_clear_state();
1902 assert!(player.in_flight.is_none(), "the stop consumed it");
1903 }
1904
1905 #[test]
1906 fn resume_with_nothing_loaded_plays_the_cursor() {
1907 let dir = tempfile::tempdir().unwrap();
1908 let path = dir.path().join("t.wav");
1909 crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
1910
1911 let mut player = Player::new();
1912 player.backend = Box::new(StuckBackend {
1913 rate: 8_000.0,
1914 asked: Default::default(),
1915 });
1916 let item = PlaylistItem {
1917 db_id: Some(5),
1918 path,
1919 ..make_item("t")
1920 };
1921 let id = item.id;
1922 player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
1923 player.shared_state.set_cursor(Some(id));
1924 assert!(player.active_playback.is_none());
1925
1926 player.process_command(PlayerCommand::Resume);
1927 assert!(
1928 player.active_playback.is_some(),
1929 "the cursor's track started"
1930 );
1931 assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
1932 player.process_command(PlayerCommand::Stop);
1933 }
1934
1935 #[test]
1936 fn the_server_hears_each_turn_playback_takes() {
1937 use PlaybackReportState::{Paused, Playing, Stopped};
1938 use history::PlaybackReport;
1939
1940 let dir = tempfile::tempdir().unwrap();
1941 let path = dir.path().join("t.wav");
1942 crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
1943
1944 let mut player = Player::new();
1945 player.backend = Box::new(StuckBackend {
1946 rate: 8_000.0,
1947 asked: Default::default(),
1948 });
1949 let (recorder, events) = PlayRecorder::capture();
1950 player.history = Some(recorder);
1951
1952 let item = PlaylistItem {
1953 db_id: Some(5),
1954 path,
1955 ..make_item("t")
1956 };
1957 let id = item.id;
1958 player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
1959 player.process_command(PlayerCommand::Play(id));
1960 player.process_command(PlayerCommand::Pause);
1961 player.process_command(PlayerCommand::Seek(4_000));
1962 player.process_command(PlayerCommand::Resume);
1963 player.process_command(PlayerCommand::Stop);
1964
1965 let report = |state, position_ms| {
1966 PlayEvent::Playback(PlaybackReport {
1967 track_id: 5,
1968 state,
1969 position_ms,
1970 })
1971 };
1972 assert_eq!(
1973 events.try_iter().collect::<Vec<_>>(),
1974 vec![
1975 PlayEvent::Started {
1976 track_id: 5,
1977 position_ms: 0
1978 },
1979 report(Paused, 0),
1980 report(Paused, 4_000),
1981 report(Playing, 4_000),
1982 report(Stopped, 4_000),
1983 PlayEvent::Finished {
1984 track_id: 5,
1985 listened_ms: 0
1986 },
1987 ]
1988 );
1989 }
1990
1991 #[test]
1992 fn removing_the_playing_track_resumes_at_its_successor() {
1993 let mut player = Player::new();
1994 let ids = seed(&mut player, 5);
1995 player.shared_state.set_cursor(Some(ids[2]));
1996
1997 player.process_command(PlayerCommand::RemoveFromPlaylist(ids[2]));
1998
1999 assert_eq!(
2000 player.shared_state.cursor(),
2001 Some(ids[3]),
2002 "playback must continue at the next track, not restart the queue"
2003 );
2004 assert_eq!(player.playback_starts, 1);
2005 }
2006
2007 #[test]
2008 fn removing_the_first_playing_track_resumes_at_the_new_first() {
2009 let mut player = Player::new();
2010 let ids = seed(&mut player, 3);
2011 player.shared_state.set_cursor(Some(ids[0]));
2012
2013 player.process_command(PlayerCommand::RemoveFromPlaylist(ids[0]));
2014
2015 assert_eq!(player.shared_state.cursor(), Some(ids[1]));
2016 }
2017
2018 #[test]
2019 fn next_track_parks_on_a_track_that_has_not_downloaded_yet() {
2020 let mut player = Player::new();
2021 let playing = make_item("playing");
2022 let waiting = pending_item("waiting");
2023 let later = make_item("later");
2024 let (playing_id, waiting_id) = (playing.id, waiting.id);
2025 player.process_command(PlayerCommand::AddToPlaylist(vec![playing, waiting, later]));
2026 player.shared_state.set_cursor(Some(playing_id));
2027
2028 player.process_command(PlayerCommand::DecodeFinished);
2029
2030 assert_eq!(
2031 player.shared_state.cursor(),
2032 Some(waiting_id),
2033 "the cursor parks on the track being fetched"
2034 );
2035 assert_eq!(
2036 player.playback_starts, 0,
2037 "nothing to play until its bytes land"
2038 );
2039
2040 player
2043 .shared_state
2044 .update_item_state(waiting_id, ItemState::Ready);
2045 player.process_command(PlayerCommand::TrackReady(waiting_id));
2046
2047 assert_eq!(player.playback_starts, 1);
2048 assert_eq!(player.shared_state.cursor(), Some(waiting_id));
2049 }
2050
2051 #[test]
2052 fn a_download_that_cannot_land_moves_the_cursor_on() {
2053 let mut player = Player::new();
2054 let waiting = pending_item("waiting");
2055 let later = make_item("later");
2056 let (waiting_id, later_id) = (waiting.id, later.id);
2057 player.process_command(PlayerCommand::AddToPlaylist(vec![waiting, later]));
2058
2059 player.process_command(PlayerCommand::Play(waiting_id));
2060 assert_eq!(player.playback_starts, 0, "nothing to play yet");
2061
2062 player
2064 .shared_state
2065 .update_item_state(waiting_id, ItemState::Failed("remote unavailable".into()));
2066 player.process_command(PlayerCommand::TrackFailed(waiting_id));
2067
2068 assert_eq!(
2069 player.shared_state.cursor(),
2070 Some(later_id),
2071 "the queue moves past a track that can never load"
2072 );
2073 assert_eq!(player.playback_starts, 1);
2074 }
2075
2076 #[test]
2077 fn a_queue_that_can_never_load_stops_rather_than_waiting() {
2078 let mut player = Player::new();
2079 let first = pending_item("first");
2080 let second = pending_item("second");
2081 let (first_id, second_id) = (first.id, second.id);
2082 player.process_command(PlayerCommand::AddToPlaylist(vec![first, second]));
2083
2084 player.process_command(PlayerCommand::Play(first_id));
2085 for id in [first_id, second_id] {
2086 player
2087 .shared_state
2088 .update_item_state(id, ItemState::Failed("remote unavailable".into()));
2089 player.process_command(PlayerCommand::TrackFailed(id));
2090 }
2091
2092 assert_eq!(player.playback_starts, 0);
2093 assert_eq!(
2094 player.shared_state.playback_state(),
2095 PlaybackState::Stopped,
2096 "a stop the UI can see, not an indefinite wait for TrackReady"
2097 );
2098 }
2099
2100 #[test]
2101 fn a_failure_elsewhere_in_the_queue_leaves_the_cursor_alone() {
2102 let mut player = Player::new();
2103 let waiting = pending_item("waiting");
2104 let other = pending_item("other");
2105 let (waiting_id, other_id) = (waiting.id, other.id);
2106 player.process_command(PlayerCommand::AddToPlaylist(vec![waiting, other]));
2107 player.process_command(PlayerCommand::Play(waiting_id));
2108
2109 player
2110 .shared_state
2111 .update_item_state(other_id, ItemState::Failed("remote unavailable".into()));
2112 player.process_command(PlayerCommand::TrackFailed(other_id));
2113
2114 assert_eq!(
2115 player.shared_state.cursor(),
2116 Some(waiting_id),
2117 "a track still downloading keeps the cursor"
2118 );
2119 }
2120
2121 #[test]
2122 fn batch_delete_containing_the_cursor_restarts_the_engine_once() {
2123 let mut player = Player::new();
2124 let ids = seed(&mut player, 5);
2125 player.shared_state.set_cursor(Some(ids[2]));
2126
2127 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![
2128 ids[1], ids[2], ids[3],
2129 ]));
2130
2131 assert_eq!(playlist_titles(&player), vec!["t0", "t4"]);
2132 assert_eq!(player.shared_state.cursor(), Some(ids[4]));
2133 assert_eq!(
2134 player.playback_starts, 1,
2135 "one resume for the whole selection, not one per deleted track"
2136 );
2137 }
2138
2139 #[test]
2140 fn batch_delete_below_the_cursor_leaves_playback_alone() {
2141 let mut player = Player::new();
2142 let ids = seed(&mut player, 4);
2143 player.shared_state.set_cursor(Some(ids[0]));
2144
2145 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![ids[2], ids[3]]));
2146
2147 assert_eq!(player.shared_state.cursor(), Some(ids[0]));
2148 assert_eq!(player.playback_starts, 0);
2149 }
2150
2151 #[test]
2152 fn undo_of_a_batch_delete_restores_the_original_order() {
2153 let mut player = Player::new();
2157 let items = vec![
2158 make_item("A"),
2159 make_item("B"),
2160 make_item("C"),
2161 make_item("D"),
2162 ];
2163 let (b_id, c_id) = (items[1].id, items[2].id);
2164 player.process_command(PlayerCommand::AddToPlaylist(items));
2165
2166 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![c_id, b_id]));
2167 assert_eq!(playlist_titles(&player), vec!["A", "D"]);
2168
2169 player.process_command(PlayerCommand::Undo);
2170 assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
2171 }
2172
2173 #[test]
2176 fn undo_add_removes_items() {
2177 let mut player = Player::new();
2178 let items = vec![make_item("A"), make_item("B")];
2179 let ids: Vec<_> = items.iter().map(|i| i.id).collect();
2180
2181 player.process_command(PlayerCommand::AddToPlaylist(items));
2182 assert_eq!(playlist_ids(&player), ids);
2183 assert!(player.undo_stack().can_undo());
2184
2185 player.process_command(PlayerCommand::Undo);
2186 assert!(playlist_ids(&player).is_empty());
2187 assert!(player.undo_stack().can_redo());
2188 }
2189
2190 #[test]
2191 fn redo_add_restores_items() {
2192 let mut player = Player::new();
2193 let items = vec![make_item("A"), make_item("B")];
2194
2195 player.process_command(PlayerCommand::AddToPlaylist(items));
2196 player.process_command(PlayerCommand::Undo);
2197 assert!(playlist_ids(&player).is_empty());
2198
2199 player.process_command(PlayerCommand::Redo);
2200 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2201 }
2202
2203 #[test]
2206 fn undo_remove_restores_item_at_position() {
2207 let mut player = Player::new();
2208 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2209 let b_id = items[1].id;
2210
2211 player.process_command(PlayerCommand::AddToPlaylist(items));
2212 player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
2213 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2214
2215 player.process_command(PlayerCommand::Undo);
2216 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2217 }
2218
2219 #[test]
2220 fn undo_remove_first_item() {
2221 let mut player = Player::new();
2222 let items = vec![make_item("A"), make_item("B")];
2223 let a_id = items[0].id;
2224
2225 player.process_command(PlayerCommand::AddToPlaylist(items));
2226 player.process_command(PlayerCommand::RemoveFromPlaylist(a_id));
2227 assert_eq!(playlist_titles(&player), vec!["B"]);
2228
2229 player.process_command(PlayerCommand::Undo);
2230 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2231 }
2232
2233 #[test]
2234 fn undo_batch_remove_restores_all() {
2235 let mut player = Player::new();
2236 let items = vec![
2237 make_item("A"),
2238 make_item("B"),
2239 make_item("C"),
2240 make_item("D"),
2241 ];
2242 let b_id = items[1].id;
2243 let c_id = items[2].id;
2244
2245 player.process_command(PlayerCommand::AddToPlaylist(items));
2246 let version_before = player.shared_state.playlist_version();
2247 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![b_id, c_id]));
2248 assert_eq!(playlist_titles(&player), vec!["A", "D"]);
2249 assert_eq!(
2252 player.shared_state.playlist_version(),
2253 version_before + 1,
2254 "batch removal must bump the playlist version exactly once"
2255 );
2256
2257 player.process_command(PlayerCommand::Undo);
2259 assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
2260 }
2261
2262 #[test]
2263 fn redo_batch_remove() {
2264 let mut player = Player::new();
2265 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2266 let a_id = items[0].id;
2267 let b_id = items[1].id;
2268
2269 player.process_command(PlayerCommand::AddToPlaylist(items));
2270 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![a_id, b_id]));
2271 player.process_command(PlayerCommand::Undo);
2272 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2273
2274 player.process_command(PlayerCommand::Redo);
2275 assert_eq!(playlist_titles(&player), vec!["C"]);
2276 }
2277
2278 #[test]
2279 fn redo_remove() {
2280 let mut player = Player::new();
2281 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2282 let b_id = items[1].id;
2283
2284 player.process_command(PlayerCommand::AddToPlaylist(items));
2285 player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
2286 player.process_command(PlayerCommand::Undo);
2287 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2288
2289 player.process_command(PlayerCommand::Redo);
2290 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2291 }
2292
2293 #[test]
2296 fn undo_insert_removes_inserted_items() {
2297 let mut player = Player::new();
2298 let items = vec![make_item("A"), make_item("C")];
2299 let a_id = items[0].id;
2300
2301 player.process_command(PlayerCommand::AddToPlaylist(items));
2302
2303 let inserted = vec![make_item("B")];
2304 player.process_command(PlayerCommand::InsertInPlaylist {
2305 items: inserted,
2306 after: a_id,
2307 });
2308 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2309
2310 player.process_command(PlayerCommand::Undo);
2311 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2312 }
2313
2314 #[test]
2317 fn undo_move_restores_position() {
2318 let mut player = Player::new();
2319 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2320 let a_id = items[0].id;
2321 let c_id = items[2].id;
2322
2323 player.process_command(PlayerCommand::AddToPlaylist(items));
2324
2325 player.process_command(PlayerCommand::MoveInPlaylist {
2327 id: a_id,
2328 target: c_id,
2329 after: true,
2330 });
2331 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2332
2333 player.process_command(PlayerCommand::Undo);
2334 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2335 }
2336
2337 #[test]
2338 fn redo_move() {
2339 let mut player = Player::new();
2340 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2341 let a_id = items[0].id;
2342 let c_id = items[2].id;
2343
2344 player.process_command(PlayerCommand::AddToPlaylist(items));
2345 player.process_command(PlayerCommand::MoveInPlaylist {
2346 id: a_id,
2347 target: c_id,
2348 after: true,
2349 });
2350 player.process_command(PlayerCommand::Undo);
2351 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2352
2353 player.process_command(PlayerCommand::Redo);
2354 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2355 }
2356
2357 #[test]
2360 fn undo_batch_move() {
2361 let mut player = Player::new();
2362 let items = vec![
2363 make_item("A"),
2364 make_item("B"),
2365 make_item("C"),
2366 make_item("D"),
2367 ];
2368 let a_id = items[0].id;
2369 let b_id = items[1].id;
2370 let d_id = items[3].id;
2371
2372 player.process_command(PlayerCommand::AddToPlaylist(items));
2373
2374 player.process_command(PlayerCommand::MoveItemsInPlaylist {
2376 ids: vec![a_id, b_id],
2377 target: d_id,
2378 after: true,
2379 });
2380 assert_eq!(playlist_titles(&player), vec!["C", "D", "A", "B"]);
2381
2382 player.process_command(PlayerCommand::Undo);
2383 assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
2384 }
2385
2386 #[test]
2389 fn undo_clear_restores_playlist() {
2390 let mut player = Player::new();
2391 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2392
2393 player.process_command(PlayerCommand::AddToPlaylist(items));
2394 player.process_command(PlayerCommand::ClearPlaylist);
2395 assert!(playlist_ids(&player).is_empty());
2396
2397 player.process_command(PlayerCommand::Undo);
2398 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2399 }
2400
2401 #[test]
2406 fn undoing_a_replace_does_not_leave_the_engine_on_an_orphaned_track() {
2407 let mut player = Player::new();
2408 let original = seed(&mut player, 3);
2409 player.shared_state.set_cursor(Some(original[0]));
2410 pretend_playing(&mut player, original[0]);
2411
2412 let replacement = vec![make_item("something else")];
2413 let orphan = replacement[0].id;
2414 player.process_command(PlayerCommand::ReplacePlaylist {
2415 items: replacement,
2416 start: 0,
2417 });
2418 pretend_playing(&mut player, orphan);
2420
2421 player.process_command(PlayerCommand::Undo);
2422
2423 assert_eq!(playlist_ids(&player), original, "the queue comes back");
2424 assert!(
2425 player.shared_state.get_item(orphan).is_none(),
2426 "and the replacement is gone from it"
2427 );
2428 assert!(
2429 playing_id(&player).is_none_or(|id| player.shared_state.get_item(id).is_some()),
2430 "so nothing may still be playing out of it"
2431 );
2432 }
2433
2434 #[test]
2436 fn undoing_an_add_does_not_leave_the_engine_on_a_removed_track() {
2437 let mut player = Player::new();
2438 seed(&mut player, 2);
2439 let added = seed(&mut player, 1);
2440 pretend_playing(&mut player, added[0]);
2441
2442 player.process_command(PlayerCommand::Undo);
2443
2444 assert!(player.shared_state.get_item(added[0]).is_none());
2445 assert!(
2446 playing_id(&player).is_none_or(|id| player.shared_state.get_item(id).is_some()),
2447 "the engine cannot be left on the item the undo removed"
2448 );
2449 }
2450
2451 #[test]
2453 fn undoing_a_move_leaves_playback_alone() {
2454 let mut player = Player::new();
2455 let ids = seed(&mut player, 3);
2456 player.shared_state.set_cursor(Some(ids[0]));
2457 pretend_playing(&mut player, ids[0]);
2458 let starts = player.playback_starts;
2459
2460 player.process_command(PlayerCommand::MoveInPlaylist {
2461 id: ids[2],
2462 target: ids[0],
2463 after: false,
2464 });
2465 player.process_command(PlayerCommand::Undo);
2466
2467 assert_eq!(playlist_ids(&player), ids);
2468 assert_eq!(playing_id(&player), Some(ids[0]), "still on the same track");
2469 assert_eq!(player.playback_starts, starts, "and not restarted");
2470 }
2471
2472 #[test]
2473 fn redo_clear() {
2474 let mut player = Player::new();
2475 let items = vec![make_item("A"), make_item("B")];
2476
2477 player.process_command(PlayerCommand::AddToPlaylist(items));
2478 player.process_command(PlayerCommand::ClearPlaylist);
2479 player.process_command(PlayerCommand::Undo);
2480 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2481
2482 player.process_command(PlayerCommand::Redo);
2483 assert!(playlist_ids(&player).is_empty());
2484 }
2485
2486 #[test]
2489 fn multiple_undos_in_sequence() {
2490 let mut player = Player::new();
2491
2492 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("A")]));
2493 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
2494 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("C")]));
2495 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2496
2497 player.process_command(PlayerCommand::Undo);
2498 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2499
2500 player.process_command(PlayerCommand::Undo);
2501 assert_eq!(playlist_titles(&player), vec!["A"]);
2502
2503 player.process_command(PlayerCommand::Undo);
2504 assert!(playlist_ids(&player).is_empty());
2505 }
2506
2507 #[test]
2508 fn undo_redo_undo_cycle() {
2509 let mut player = Player::new();
2510 let items = vec![make_item("A"), make_item("B")];
2511
2512 player.process_command(PlayerCommand::AddToPlaylist(items));
2513 player.process_command(PlayerCommand::Undo);
2514 assert!(playlist_ids(&player).is_empty());
2515
2516 player.process_command(PlayerCommand::Redo);
2517 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2518
2519 player.process_command(PlayerCommand::Undo);
2520 assert!(playlist_ids(&player).is_empty());
2521 }
2522
2523 #[test]
2524 fn new_action_clears_redo_stack() {
2525 let mut player = Player::new();
2526 let items = vec![make_item("A")];
2527
2528 player.process_command(PlayerCommand::AddToPlaylist(items));
2529 player.process_command(PlayerCommand::Undo);
2530 assert!(player.undo_stack().can_redo());
2531
2532 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
2534 assert!(!player.undo_stack().can_redo());
2535 }
2536
2537 #[test]
2538 fn undo_on_empty_stack_is_noop() {
2539 let mut player = Player::new();
2540 player.process_command(PlayerCommand::Undo);
2541 assert!(playlist_ids(&player).is_empty());
2542 }
2543
2544 #[test]
2545 fn redo_on_empty_stack_is_noop() {
2546 let mut player = Player::new();
2547 player.process_command(PlayerCommand::Redo);
2548 assert!(playlist_ids(&player).is_empty());
2549 }
2550
2551 #[test]
2554 fn playback_commands_not_undoable() {
2555 let mut player = Player::new();
2556 player.process_command(PlayerCommand::Pause);
2557 player.process_command(PlayerCommand::Resume);
2558 player.process_command(PlayerCommand::NextTrack);
2559 player.process_command(PlayerCommand::PrevTrack);
2560 assert!(!player.undo_stack().can_undo());
2561 }
2562
2563 #[test]
2564 fn update_paths_not_undoable() {
2565 let mut player = Player::new();
2566 let items = vec![make_item("A")];
2567 let id = items[0].id;
2568 player.process_command(PlayerCommand::AddToPlaylist(items));
2569
2570 let undo_count = player.undo_stack().undo_len();
2571 player.process_command(PlayerCommand::UpdatePaths(vec![(
2572 id,
2573 PathBuf::from("/new/path.flac"),
2574 )]));
2575 assert_eq!(player.undo_stack().undo_len(), undo_count);
2576 }
2577
2578 #[test]
2581 fn add_remove_undo_undo_produces_original() {
2582 let mut player = Player::new();
2583 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2584 let b_id = items[1].id;
2585 let original_titles = vec!["A", "B", "C"];
2586
2587 player.process_command(PlayerCommand::AddToPlaylist(items));
2588 player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
2589 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2590
2591 player.process_command(PlayerCommand::Undo);
2593 assert_eq!(playlist_titles(&player), original_titles);
2594
2595 player.process_command(PlayerCommand::Undo);
2597 assert!(playlist_ids(&player).is_empty());
2598 }
2599
2600 #[test]
2601 fn interleaved_adds_and_moves_undo() {
2602 let mut player = Player::new();
2603 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2604 let a_id = items[0].id;
2605 let c_id = items[2].id;
2606
2607 player.process_command(PlayerCommand::AddToPlaylist(items));
2608
2609 player.process_command(PlayerCommand::MoveInPlaylist {
2611 id: a_id,
2612 target: c_id,
2613 after: true,
2614 });
2615 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2616
2617 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("D")]));
2619 assert_eq!(playlist_titles(&player), vec!["B", "C", "A", "D"]);
2620
2621 player.process_command(PlayerCommand::Undo);
2623 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2624
2625 player.process_command(PlayerCommand::Undo);
2627 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2628 }
2629
2630 #[test]
2635 fn stop_engine_drops_engine_synchronously() {
2636 use std::sync::atomic::{AtomicBool, Ordering};
2637
2638 struct MockEngine {
2639 dropped: Arc<AtomicBool>,
2640 }
2641 impl AudioEngineHandle for MockEngine {
2642 fn start(&self) -> Result<(), BackendError> {
2643 Ok(())
2644 }
2645 fn stop(&self) -> Result<(), BackendError> {
2646 Ok(())
2647 }
2648 fn is_running(&self) -> bool {
2649 false
2650 }
2651 fn fade_out(&self) {}
2652 fn fade_in(&self) -> Result<(), BackendError> {
2653 Ok(())
2654 }
2655 fn is_silent(&self) -> bool {
2656 false
2657 }
2658 }
2659 impl Drop for MockEngine {
2660 fn drop(&mut self) {
2661 self.dropped.store(true, Ordering::SeqCst);
2662 }
2663 }
2664
2665 let dropped = Arc::new(AtomicBool::new(false));
2666
2667 let stop_flag = Arc::new(AtomicBool::new(false));
2669 let decode_handle = buffer::DecodeHandle::new_for_test(stop_flag);
2670
2671 let mut player = Player::new();
2672 player.active_playback = Some(ActivePlayback {
2673 engine: Box::new(MockEngine {
2674 dropped: dropped.clone(),
2675 }),
2676 decode_handle,
2677 stream: None,
2678 _rate_watch: None,
2679 });
2680
2681 player.stop_engine();
2682
2683 assert!(
2687 dropped.load(Ordering::SeqCst),
2688 "AudioEngine must be dropped synchronously in stop_engine (GitHub #89)"
2689 );
2690 }
2691
2692 #[test]
2693 fn abandoning_a_stream_wakes_a_reader_parked_at_the_write_head() {
2694 let live = LiveStream {
2695 feed: crate::remote::downloads::ByteFeed::new(),
2696 abandoned: Default::default(),
2697 };
2698 let feed = live.feed.clone();
2699 let started = std::time::Instant::now();
2700 let reader = thread::spawn(move || {
2701 feed.wait_past(
2702 0,
2703 std::time::Instant::now() + std::time::Duration::from_secs(30),
2704 )
2705 });
2706 thread::sleep(std::time::Duration::from_millis(50));
2707 live.abandon();
2708 reader.join().unwrap();
2709
2710 assert!(live.abandoned.load(std::sync::atomic::Ordering::Acquire));
2711 assert!(started.elapsed() < std::time::Duration::from_secs(5));
2712 }
2713
2714 struct StuckBackend {
2719 rate: f64,
2720 asked: Arc<std::sync::Mutex<Option<(f64, u32)>>>,
2721 }
2722
2723 struct NullEngine;
2724 impl AudioEngineHandle for NullEngine {
2725 fn start(&self) -> Result<(), BackendError> {
2726 Ok(())
2727 }
2728 fn stop(&self) -> Result<(), BackendError> {
2729 Ok(())
2730 }
2731 fn is_running(&self) -> bool {
2732 false
2733 }
2734 fn fade_out(&self) {}
2735 fn fade_in(&self) -> Result<(), BackendError> {
2736 Ok(())
2737 }
2738 fn is_silent(&self) -> bool {
2739 false
2740 }
2741 }
2742
2743 impl AudioBackend for StuckBackend {
2744 fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
2745 Ok(vec![self.default_device()?])
2746 }
2747 fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
2748 Ok(backend::DeviceInfo {
2749 name: "Stuck DAC".into(),
2750 sample_rates: vec![self.rate],
2751 platform_id: 0,
2752 })
2753 }
2754 fn supported_sample_rates(
2755 &self,
2756 _device: &backend::DeviceInfo,
2757 ) -> Result<Vec<f64>, BackendError> {
2758 Ok(vec![self.rate])
2759 }
2760 fn get_device_sample_rate(
2761 &self,
2762 _device: &backend::DeviceInfo,
2763 ) -> Result<f64, BackendError> {
2764 Ok(self.rate)
2765 }
2766 fn set_device_sample_rate(
2767 &self,
2768 _device: &backend::DeviceInfo,
2769 rate: f64,
2770 ) -> Result<f64, BackendError> {
2771 Err(BackendError::UnsupportedSampleRate(rate))
2772 }
2773 fn create_engine(
2774 &self,
2775 _device: &backend::DeviceInfo,
2776 sample_rate: f64,
2777 channels: u32,
2778 _consumer: rtrb::Consumer<f32>,
2779 _samples_played: Arc<AtomicU64>,
2780 ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
2781 *self.asked.lock().unwrap() = Some((sample_rate, channels));
2782 Ok(Box::new(NullEngine))
2783 }
2784 }
2785
2786 fn engine_format_for(source_rate: u32, channels: u16, device_rate: f64) -> (f64, u32) {
2787 let asked = Arc::new(std::sync::Mutex::new(None));
2788 let mut player = Player::new();
2789 player.backend = Box::new(StuckBackend {
2790 rate: device_rate,
2791 asked: asked.clone(),
2792 });
2793
2794 let info = buffer::StreamInfo {
2795 codec: "MP3".into(),
2796 sample_rate: source_rate,
2797 channels,
2798 bit_depth: Some(16),
2799 bitrate_kbps: None,
2800 duration_ms: 1000,
2801 };
2802 let (_producer, consumer) = rtrb::RingBuffer::new(16);
2803 player
2804 .create_engine_for(&info, consumer)
2805 .expect("engine creation should succeed");
2806 let asked = *asked.lock().unwrap();
2807 asked.expect("engine was never created")
2808 }
2809
2810 #[test]
2811 fn engine_uses_source_rate_when_device_refuses_switch() {
2812 assert_eq!(engine_format_for(22050, 2, 48000.0), (22050.0, 2));
2815 assert_eq!(engine_format_for(32000, 2, 44100.0), (32000.0, 2));
2816 }
2817
2818 #[test]
2819 fn engine_uses_source_channel_count() {
2820 assert_eq!(engine_format_for(44100, 1, 44100.0), (44100.0, 1));
2821 }
2822
2823 fn settled_rate_for(source_rate: u32, device_rate: f64) -> Option<u32> {
2825 let mut player = Player::new();
2826 player.backend = Box::new(StuckBackend {
2827 rate: device_rate,
2828 asked: Arc::new(std::sync::Mutex::new(None)),
2829 });
2830 let state = player.shared_state.clone();
2831
2832 let info = buffer::StreamInfo {
2833 codec: "MP3".into(),
2834 sample_rate: source_rate,
2835 channels: 2,
2836 bit_depth: Some(16),
2837 bitrate_kbps: None,
2838 duration_ms: 1000,
2839 };
2840 let (_producer, consumer) = rtrb::RingBuffer::new(16);
2841 player
2842 .create_engine_for(&info, consumer)
2843 .expect("engine creation should succeed");
2844 state.output_sample_rate()
2845 }
2846
2847 #[test]
2848 fn settled_device_rate_reaches_the_shared_state() {
2849 assert_eq!(settled_rate_for(22050, 48000.0), Some(48000));
2853 assert_eq!(settled_rate_for(44100, 44100.0), Some(44100));
2855 }
2856
2857 struct SlowBackend {
2859 observed: Arc<std::sync::Mutex<Vec<Option<u32>>>>,
2860 state: Arc<SharedPlayerState>,
2861 }
2862
2863 impl AudioBackend for SlowBackend {
2864 fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
2865 Ok(vec![self.default_device()?])
2866 }
2867 fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
2868 Ok(backend::DeviceInfo {
2869 name: "Slow DAC".into(),
2870 sample_rates: vec![44100.0, 48000.0],
2871 platform_id: 0,
2872 })
2873 }
2874 fn supported_sample_rates(
2875 &self,
2876 _device: &backend::DeviceInfo,
2877 ) -> Result<Vec<f64>, BackendError> {
2878 Ok(vec![44100.0, 48000.0])
2879 }
2880 fn get_device_sample_rate(
2881 &self,
2882 _device: &backend::DeviceInfo,
2883 ) -> Result<f64, BackendError> {
2884 Ok(48000.0)
2885 }
2886 fn set_device_sample_rate(
2887 &self,
2888 _device: &backend::DeviceInfo,
2889 rate: f64,
2890 ) -> Result<f64, BackendError> {
2891 self.observed
2893 .lock()
2894 .unwrap()
2895 .push(self.state.output_sample_rate());
2896 Ok(rate)
2897 }
2898 fn create_engine(
2899 &self,
2900 _device: &backend::DeviceInfo,
2901 _sample_rate: f64,
2902 _channels: u32,
2903 _consumer: rtrb::Consumer<f32>,
2904 _samples_played: Arc<AtomicU64>,
2905 ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
2906 Ok(Box::new(NullEngine))
2907 }
2908 }
2909
2910 #[test]
2911 fn the_previous_rate_is_not_published_while_the_device_reclocks() {
2912 let mut player = Player::new();
2918 let state = player.shared_state.clone();
2919 state.set_output_sample_rate(48000);
2920
2921 let observed = Arc::new(std::sync::Mutex::new(Vec::new()));
2922 player.backend = Box::new(SlowBackend {
2923 observed: observed.clone(),
2924 state: state.clone(),
2925 });
2926
2927 let info = buffer::StreamInfo {
2928 codec: "FLAC".into(),
2929 sample_rate: 44100,
2930 channels: 2,
2931 bit_depth: Some(16),
2932 bitrate_kbps: None,
2933 duration_ms: 1000,
2934 };
2935 let (_producer, consumer) = rtrb::RingBuffer::new(16);
2936 player
2937 .create_engine_for(&info, consumer)
2938 .expect("engine creation should succeed");
2939
2940 assert_eq!(
2941 *observed.lock().unwrap(),
2942 vec![None],
2943 "mid-switch the output rate must read as unknown, not as the last track's"
2944 );
2945 assert_eq!(state.output_sample_rate(), Some(44100));
2946 }
2947
2948 struct WatchedBackend {
2950 inner: StuckBackend,
2951 #[allow(clippy::type_complexity)]
2952 captured: Arc<std::sync::Mutex<Option<Box<dyn Fn(f64) + Send + Sync>>>>,
2953 }
2954
2955 struct NullWatch;
2956 impl backend::SampleRateWatch for NullWatch {}
2957
2958 impl AudioBackend for WatchedBackend {
2959 fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
2960 self.inner.list_devices()
2961 }
2962 fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
2963 self.inner.default_device()
2964 }
2965 fn supported_sample_rates(
2966 &self,
2967 device: &backend::DeviceInfo,
2968 ) -> Result<Vec<f64>, BackendError> {
2969 self.inner.supported_sample_rates(device)
2970 }
2971 fn get_device_sample_rate(
2972 &self,
2973 device: &backend::DeviceInfo,
2974 ) -> Result<f64, BackendError> {
2975 self.inner.get_device_sample_rate(device)
2976 }
2977 fn set_device_sample_rate(
2978 &self,
2979 device: &backend::DeviceInfo,
2980 rate: f64,
2981 ) -> Result<f64, BackendError> {
2982 self.inner.set_device_sample_rate(device, rate)
2983 }
2984 fn watch_device_sample_rate(
2985 &self,
2986 _device: &backend::DeviceInfo,
2987 on_change: Box<dyn Fn(f64) + Send + Sync>,
2988 ) -> Option<Box<dyn backend::SampleRateWatch>> {
2989 *self.captured.lock().unwrap() = Some(on_change);
2990 Some(Box::new(NullWatch))
2991 }
2992 fn create_engine(
2993 &self,
2994 device: &backend::DeviceInfo,
2995 sample_rate: f64,
2996 channels: u32,
2997 consumer: rtrb::Consumer<f32>,
2998 samples_played: Arc<AtomicU64>,
2999 ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
3000 self.inner
3001 .create_engine(device, sample_rate, channels, consumer, samples_played)
3002 }
3003 }
3004
3005 #[test]
3006 fn external_rate_change_reaches_the_shared_state() {
3007 let captured = Arc::new(std::sync::Mutex::new(None));
3011 let mut player = Player::new();
3012 player.backend = Box::new(WatchedBackend {
3013 inner: StuckBackend {
3014 rate: 44100.0,
3015 asked: Arc::new(std::sync::Mutex::new(None)),
3016 },
3017 captured: captured.clone(),
3018 });
3019 let state = player.shared_state.clone();
3020
3021 let info = buffer::StreamInfo {
3022 codec: "FLAC".into(),
3023 sample_rate: 44100,
3024 channels: 2,
3025 bit_depth: Some(16),
3026 bitrate_kbps: None,
3027 duration_ms: 1000,
3028 };
3029 let (_producer, consumer) = rtrb::RingBuffer::new(16);
3030 player
3031 .create_engine_for(&info, consumer)
3032 .expect("engine creation should succeed");
3033 assert_eq!(state.output_sample_rate(), Some(44100));
3034
3035 let on_change = captured.lock().unwrap().take().expect("watch registered");
3036 on_change(48000.0);
3037 assert_eq!(state.output_sample_rate(), Some(48000));
3038 }
3039}