1pub mod commands;
2pub mod history;
3pub mod state;
4pub mod undo;
5
6use std::path::{Path, PathBuf};
7use std::sync::Arc;
8use std::thread;
9
10use thiserror::Error;
11
12use crate::audio::{
13 analyzer::VizAnalyzer,
14 backend::{self, AudioBackend, AudioEngineHandle, BackendError, SampleRateWatch},
15 buffer, streaming,
16 viz::{VizBuffer, VizSnapshot},
17};
18use buffer::PlaybackTimeline;
19use commands::{CommandChannel, PlayerCommand};
20use history::{InFlight, PlayEvent, PlayRecorder};
21use state::{
22 ItemState, LoadState, PlaybackSource, PlaybackState, QueueItemId, SharedPlayerState, TrackInfo,
23};
24use undo::{UndoEntry, UndoStack};
25
26pub(crate) const RING_BUFFER_SIZE: usize = 192_000 * 2;
28
29const SEEK_END_GUARD_MS: u64 = 500;
32
33#[derive(Debug, Error)]
34pub enum PlayerError {
35 #[error("backend error: {0}")]
36 Backend(#[from] BackendError),
37 #[error("decode error: {0}")]
38 Decode(#[from] buffer::DecodeError),
39}
40
41#[derive(Clone)]
44struct StreamSource {
45 path: PathBuf,
46 bytes_written: Arc<crate::remote::downloads::ByteFeed>,
47 total: u64,
48 mode: streaming::ProbeMode,
49}
50
51fn hint_for(path: &Path) -> symphonia::core::formats::probe::Hint {
53 let mut hint = symphonia::core::formats::probe::Hint::new();
54 if let Some(ext) = path.extension().and_then(|e| e.to_str()) {
55 hint.with_extension(ext);
56 }
57 hint
58}
59
60fn lengthless_mode_for(path: &Path) -> streaming::ProbeMode {
68 let ogg = path
69 .extension()
70 .and_then(|e| e.to_str())
71 .is_some_and(|ext| {
72 matches!(
73 ext.to_ascii_lowercase().as_str(),
74 "ogg" | "oga" | "opus" | "spx"
75 )
76 });
77 if ogg {
78 streaming::ProbeMode::LengthlessWholeEnd
79 } else {
80 streaming::ProbeMode::Lengthless
81 }
82}
83
84pub struct Player {
86 shared_state: Arc<SharedPlayerState>,
87 commands: CommandChannel,
88 active_playback: Option<ActivePlayback>,
89 timeline: Arc<PlaybackTimeline>,
90 viz_buffer: Arc<VizBuffer>,
91 viz_snapshot: Arc<VizSnapshot>,
92 _viz_analyzer: VizAnalyzer,
94 undo_stack: UndoStack,
95 batch_buffer: Option<Vec<UndoEntry>>,
98 output_device_name: Option<String>,
100 backend: Box<dyn AudioBackend>,
102 last_skip: std::time::Instant,
104 stream_mode: streaming::ProbeMode,
107 history: Option<PlayRecorder>,
110 in_flight: Option<InFlight>,
112 #[cfg(test)]
115 playback_starts: usize,
116}
117
118struct ActivePlayback {
120 engine: Box<dyn AudioEngineHandle>,
121 decode_handle: buffer::DecodeHandle,
122 stream: Option<LiveStream>,
124 _rate_watch: Option<Box<dyn SampleRateWatch>>,
127}
128
129struct LiveStream {
133 feed: Arc<crate::remote::downloads::ByteFeed>,
134 abandoned: Arc<std::sync::atomic::AtomicBool>,
135}
136
137impl LiveStream {
138 fn abandon(&self) {
139 self.abandoned
140 .store(true, std::sync::atomic::Ordering::Release);
141 self.feed.done();
142 }
143}
144
145impl Default for Player {
146 fn default() -> Self {
147 Self::new()
148 }
149}
150
151impl Player {
152 pub fn new() -> Self {
153 let viz_buffer = VizBuffer::new();
154 let viz_snapshot = VizSnapshot::new();
155 let timeline = PlaybackTimeline::new();
156 let cfg = crate::config::Config::load_or_default();
157 let viz_analyzer = VizAnalyzer::spawn_with_snapshot(
158 Arc::clone(&viz_buffer),
159 &cfg.visualizer,
160 Arc::clone(&viz_snapshot),
161 timeline.samples_played_counter(),
162 );
163
164 Self {
165 shared_state: SharedPlayerState::new(),
166 commands: CommandChannel::new(),
167 active_playback: None,
168 timeline,
169 viz_buffer,
170 viz_snapshot,
171 _viz_analyzer: viz_analyzer,
172 undo_stack: UndoStack::new(),
173 batch_buffer: None,
174 output_device_name: cfg.playback.output_device.clone(),
175 backend: crate::audio::platform_backend(),
176 last_skip: std::time::Instant::now(),
177 stream_mode: streaming::ProbeMode::Full,
178 history: None,
179 in_flight: None,
180 #[cfg(test)]
181 playback_starts: 0,
182 }
183 }
184
185 pub fn shared_state(&self) -> Arc<SharedPlayerState> {
187 self.shared_state.clone()
188 }
189
190 pub fn timeline(&self) -> Arc<PlaybackTimeline> {
192 self.timeline.clone()
193 }
194
195 pub fn viz_buffer(&self) -> Arc<VizBuffer> {
197 self.viz_buffer.clone()
198 }
199
200 pub fn viz_snapshot(&self) -> Arc<VizSnapshot> {
203 self.viz_snapshot.clone()
204 }
205
206 pub fn undo_stack(&self) -> &UndoStack {
208 &self.undo_stack
209 }
210
211 #[allow(clippy::type_complexity)]
219 fn create_engine_for(
220 &self,
221 info: &buffer::StreamInfo,
222 consumer: rtrb::Consumer<f32>,
223 ) -> Result<(Box<dyn AudioEngineHandle>, Option<Box<dyn SampleRateWatch>>), PlayerError> {
224 let device = self.resolve_device()?;
225 let device_rate = self.backend.get_device_sample_rate(&device)?;
226 let source_rate = info.sample_rate as f64;
227
228 self.shared_state.clear_output_sample_rate();
232
233 let settled = if (device_rate - source_rate).abs() > 0.1 {
234 log::info!(
235 "switching device sample rate: {}Hz → {}Hz",
236 device_rate,
237 source_rate
238 );
239 match self.backend.set_device_sample_rate(&device, source_rate) {
240 Ok(rate) => rate,
241 Err(e) => {
242 log::warn!("failed to set device sample rate: {}", e);
243 device_rate
244 }
245 }
246 } else {
247 device_rate
248 };
249
250 if (settled - source_rate).abs() > 0.1 {
251 log::warn!(
252 "device stayed at {}Hz (wanted {}Hz) — output is resampled, not bit-perfect",
253 settled,
254 source_rate
255 );
256 }
257
258 self.shared_state
261 .set_output_sample_rate(settled.round() as u32);
262
263 let watch_state = self.shared_state.clone();
267 let watch_name = device.name.clone();
268 let rate_watch = self.backend.watch_device_sample_rate(
269 &device,
270 Box::new(move |rate| {
271 log::info!("device sample rate changed externally: {rate}Hz on '{watch_name}'");
272 watch_state.set_output_sample_rate(rate.round() as u32);
273 }),
274 );
275
276 let engine = self.backend.create_engine(
277 &device,
278 source_rate,
279 info.channels as u32,
280 consumer,
281 self.timeline.samples_played_counter(),
282 )?;
283
284 Ok((engine, rate_watch))
285 }
286
287 fn resolve_device(&self) -> Result<backend::DeviceInfo, PlayerError> {
290 if let Some(ref name) = self.output_device_name {
291 match self.backend.list_devices() {
292 Ok(devices) => {
293 if let Some(dev) = devices.into_iter().find(|d| d.name == *name) {
294 return Ok(dev);
295 }
296 log::warn!(
297 "configured output device '{}' not found, falling back to default",
298 name,
299 );
300 }
301 Err(e) => {
302 log::warn!("failed to list devices while resolving '{}': {}", name, e);
303 }
304 }
305 }
306 Ok(self.backend.default_device()?)
307 }
308
309 pub fn set_output_device(&mut self, name: String) {
312 log::info!("switching output device to: {}", name);
313 self.output_device_name = Some(name.clone());
314
315 if let Err(e) = crate::config::Config::persist(|cfg| {
316 cfg.playback.output_device = Some(name);
317 }) {
318 log::error!("failed to save output device config: {}", e);
319 }
320
321 self.restart_on_current_track();
322 }
323
324 pub fn clear_output_device(&mut self) {
326 log::info!("reverting to system default output device");
327 self.output_device_name = None;
328
329 if let Err(e) = crate::config::Config::persist(|cfg| {
330 cfg.playback.output_device = None;
331 }) {
332 log::error!("failed to save output device config: {}", e);
333 }
334
335 self.restart_on_current_track();
336 }
337
338 fn restart_on_current_track(&mut self) {
341 if let Some(info) = self.shared_state.track_info() {
342 let position_ms = self.shared_state.position_ms();
343 if let Err(e) = self.restart_current(&info, position_ms) {
344 log::error!("failed to restart playback on device switch: {}", e);
345 }
346 }
347 }
348
349 pub fn output_device_name(&self) -> Option<&str> {
351 self.output_device_name.as_deref()
352 }
353
354 pub fn command_sender(&self) -> crossbeam_channel::Sender<PlayerCommand> {
356 self.commands.tx.clone()
357 }
358
359 pub fn play(&mut self, id: QueueItemId) {
362 self.shared_state.set_cursor(Some(id));
363
364 match self.shared_state.item_playback_source(id) {
365 Some(PlaybackSource::Ready(path)) => {
366 if let Err(e) = self.start_playback(id, &path, 0) {
367 log::error!("play failed: {}", e);
368 }
369 }
370 Some(PlaybackSource::Streaming {
371 path,
372 bytes_written,
373 total,
374 }) => {
375 self.stop_engine();
379 self.shared_state.set_playback_state(PlaybackState::Stopped);
380 self.probe_stream_for_playback(id, &path, bytes_written, total);
381 }
382 None => {
383 self.stop_engine();
385 self.shared_state.set_playback_state(PlaybackState::Stopped);
386 log::info!("play: item {:?} not ready, waiting for TrackReady", id);
387 }
388 }
389 }
390
391 fn start_playback(
396 &mut self,
397 id: QueueItemId,
398 path: &Path,
399 seek_ms: u64,
400 ) -> Result<(), PlayerError> {
401 #[cfg(test)]
402 {
403 self.playback_starts += 1;
404 }
405 let result = self.open_playback(id, path, seek_ms);
406 if result.is_err() {
407 self.stop_playback_and_clear_state();
408 }
409 self.wake_analyzer();
410 result
411 }
412
413 fn open_playback(
414 &mut self,
415 id: QueueItemId,
416 path: &Path,
417 seek_ms: u64,
418 ) -> Result<(), PlayerError> {
419 self.stop_engine();
420
421 let info = buffer::probe_file(path)?;
422
423 self.shared_state.set_track_info(Some(TrackInfo {
427 id,
428 path: path.to_path_buf(),
429 codec: info.codec.clone(),
430 sample_rate: info.sample_rate,
431 bit_depth: info.bit_depth,
432 bitrate_kbps: info.bitrate_kbps,
433 channels: info.channels,
434 duration_ms: info.duration_ms,
435 }));
436 self.shared_state.set_position_ms(seek_ms);
437 self.on_track_changed(id);
438 log::info!(
439 "playing: {} ({:?}) — {} {}Hz/{}ch, {}ms{}",
440 path.display(),
441 id,
442 info.codec,
443 info.sample_rate,
444 info.channels,
445 info.duration_ms,
446 if seek_ms > 0 {
447 format!(" @{}ms", seek_ms)
448 } else {
449 String::new()
450 }
451 );
452
453 let (producer, consumer) = rtrb::RingBuffer::new(RING_BUFFER_SIZE);
454
455 self.timeline.reset();
457
458 let advance_state = self.shared_state.clone();
462 let decode_cursor = parking_lot::Mutex::new(Some(id));
463 let next_track = move || {
464 let current = decode_cursor.lock().take()?;
465 let next = advance_state.peek_next_ready_after(current);
466 if let Some((next_id, _)) = &next {
467 let mut guard = decode_cursor.lock();
468 *guard = Some(*next_id);
469 }
470 next
471 };
472
473 let cfg = crate::config::Config::load_or_default();
475 let rg_mode = cfg.playback.replaygain;
476 let pre_amp_db = cfg.playback.pre_amp_db;
477
478 let finish_tx = self.commands.tx.clone();
479 let (_stream_info, decode_handle) = buffer::start_decode_file(
480 id,
481 path,
482 producer,
483 seek_ms,
484 next_track,
485 self.timeline.clone(),
486 Some(self.viz_buffer.clone()),
487 rg_mode,
488 pre_amp_db,
489 move || {
490 finish_tx.send(PlayerCommand::DecodeFinished).ok();
491 },
492 )?;
493
494 let (engine, rate_watch) = self.create_engine_for(&info, consumer)?;
495 engine.start()?;
496
497 self.shared_state.set_playback_state(PlaybackState::Playing);
498
499 self.active_playback = Some(ActivePlayback {
500 engine,
501 decode_handle,
502 stream: None,
503 _rate_watch: rate_watch,
504 });
505
506 Ok(())
507 }
508
509 fn probe_stream_for_playback(
524 &self,
525 id: QueueItemId,
526 path: &Path,
527 bytes_written: Arc<crate::remote::downloads::ByteFeed>,
528 total: u64,
529 ) {
530 let path = path.to_path_buf();
531 let tx = self.commands.tx.clone();
532 let hint = hint_for(&path);
533
534 let status = {
539 let downloading = self.stream_status_fn(id);
540 let state = self.shared_state.clone();
541 Arc::new(move || {
542 if state.is_cursor(id) {
543 downloading()
544 } else {
545 streaming::StreamStatus::Failed
546 }
547 }) as Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync>
548 };
549
550 let spawned = thread::Builder::new()
551 .name("koan-stream-probe".into())
552 .spawn(move || {
553 let attempt = |mode, wait: bool| {
560 let open = if wait {
561 streaming::PartialFileSource::open(
562 &path,
563 bytes_written.clone(),
564 total,
565 status.clone(),
566 mode,
567 )
568 } else {
569 streaming::PartialFileSource::open_for_probe(
570 &path,
571 bytes_written.clone(),
572 total,
573 status.clone(),
574 mode,
575 )
576 };
577 open.map_err(buffer::DecodeError::Io).and_then(|source| {
578 let mss = symphonia::core::io::MediaSourceStream::new(
579 Box::new(source),
580 Default::default(),
581 );
582 buffer::probe_source(mss, &hint)
583 })
584 };
585
586 let info = match attempt(streaming::ProbeMode::Full, false) {
590 Ok(info) => Some((info, streaming::ProbeMode::Full)),
591 Err(e) => {
592 log::info!(
598 "stream probe: {} needs more than has arrived ({}), opening without a length",
599 path.display(),
600 e
601 );
602 let lengthless = lengthless_mode_for(&path);
603 attempt(lengthless, true)
604 .ok()
605 .map(|info| (info, lengthless))
606 }
607 };
608
609 match info {
610 Some((info, mode)) => {
611 tx.send(PlayerCommand::StreamProbed {
612 id,
613 info: Box::new(info),
614 mode,
615 })
616 .ok();
617 }
618 None => log::info!(
621 "stream probe: {} cannot start early, waiting for the download",
622 path.display()
623 ),
624 }
625 });
626
627 if let Err(e) = spawned {
628 log::warn!("stream probe: could not spawn for {:?}: {}", id, e);
629 }
630 }
631
632 fn stream_probed(
635 &mut self,
636 id: QueueItemId,
637 info: buffer::StreamInfo,
638 mode: streaming::ProbeMode,
639 ) {
640 if !self.shared_state.is_cursor(id) {
641 return; }
643 if self.shared_state.playback_state() != PlaybackState::Stopped {
644 return; }
646
647 match self.shared_state.item_playback_source(id) {
648 Some(PlaybackSource::Ready(path)) => {
650 if let Err(e) = self.start_playback(id, &path, 0) {
651 log::error!("stream probe: playback failed: {}", e);
652 }
653 }
654 Some(PlaybackSource::Streaming {
655 path,
656 bytes_written,
657 total,
658 }) => {
659 let source = StreamSource {
660 path,
661 bytes_written,
662 total,
663 mode,
664 };
665 if let Err(e) = self.start_streaming_playback(id, source, 0, info) {
666 log::error!("stream probe: streaming playback failed: {}", e);
667 }
668 }
669 None => {}
670 }
671 }
672
673 fn stream_status_fn(
677 &self,
678 id: QueueItemId,
679 ) -> Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync> {
680 let state = self.shared_state.clone();
681 Arc::new(move || match state.item_load_state(id) {
682 Some(LoadState::Ready) => streaming::StreamStatus::Complete,
683 Some(LoadState::Failed(_)) => streaming::StreamStatus::Failed,
684 _ => streaming::StreamStatus::Downloading,
685 })
686 }
687
688 fn start_streaming_playback(
698 &mut self,
699 id: QueueItemId,
700 source: StreamSource,
701 seek_ms: u64,
702 info: buffer::StreamInfo,
703 ) -> Result<(), PlayerError> {
704 let result = self.open_streaming_playback(id, source, seek_ms, info);
705 if result.is_err() {
706 self.stop_playback_and_clear_state();
707 }
708 result
709 }
710
711 fn open_streaming_playback(
712 &mut self,
713 id: QueueItemId,
714 source: StreamSource,
715 seek_ms: u64,
716 info: buffer::StreamInfo,
717 ) -> Result<(), PlayerError> {
718 self.stop_engine();
719 self.stream_mode = source.mode;
721 let path = source.path.as_path();
722
723 let live = LiveStream {
724 feed: source.bytes_written.clone(),
725 abandoned: Default::default(),
726 };
727 let status = {
728 let downloading = self.stream_status_fn(id);
729 let abandoned = live.abandoned.clone();
730 Arc::new(move || {
731 if abandoned.load(std::sync::atomic::Ordering::Acquire) {
732 streaming::StreamStatus::Failed
733 } else {
734 downloading()
735 }
736 }) as Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync>
737 };
738 let open_source = {
739 let StreamSource {
740 path,
741 bytes_written,
742 total,
743 mode,
744 } = source.clone();
745 let status = status.clone();
746 move || {
747 streaming::PartialFileSource::open(
748 &path,
749 bytes_written.clone(),
750 total,
751 status.clone(),
752 mode,
753 )
754 }
755 };
756
757 self.shared_state.set_track_info(Some(TrackInfo {
758 id,
759 path: path.to_path_buf(),
760 codec: info.codec.clone(),
761 sample_rate: info.sample_rate,
762 bit_depth: info.bit_depth,
763 bitrate_kbps: info.bitrate_kbps,
764 channels: info.channels,
765 duration_ms: info.duration_ms,
766 }));
767 self.shared_state.set_position_ms(seek_ms);
768 self.on_track_changed(id);
769 log::info!(
770 "streaming: {} ({:?}) — {} {}Hz/{}ch, {}ms{}",
771 path.display(),
772 id,
773 info.codec,
774 info.sample_rate,
775 info.channels,
776 info.duration_ms,
777 if seek_ms > 0 {
778 format!(" @{}ms", seek_ms)
779 } else {
780 String::new()
781 },
782 );
783
784 let (producer, consumer) = rtrb::RingBuffer::new(RING_BUFFER_SIZE);
785
786 self.timeline.reset();
787
788 let advance_state = self.shared_state.clone();
790 let decode_cursor = parking_lot::Mutex::new(Some(id));
791 let next_track = move || {
792 let current = decode_cursor.lock().take()?;
793 let next = advance_state.peek_next_ready_after(current);
794 if let Some((next_id, _)) = &next {
795 let mut guard = decode_cursor.lock();
796 *guard = Some(*next_id);
797 }
798 next
799 };
800
801 let first = buffer::SourceEntry {
802 id,
803 path: path.to_path_buf(),
804 hint: hint_for(path),
805 make_mss: Box::new(move || {
806 Ok(symphonia::core::io::MediaSourceStream::new(
807 Box::new(open_source()?),
808 Default::default(),
809 ))
810 }),
811 };
812
813 let cfg = crate::config::Config::load_or_default();
815 let rg_mode = cfg.playback.replaygain;
816 let pre_amp_db = cfg.playback.pre_amp_db;
817
818 let finish_tx = self.commands.tx.clone();
819 let (_stream_info, decode_handle) = buffer::start_decode(
820 first,
821 producer,
822 seek_ms,
823 move || {
824 let (next_id, next_path) = next_track()?;
825 Some(buffer::SourceEntry::from_file(next_id, next_path))
826 },
827 self.timeline.clone(),
828 Some(self.viz_buffer.clone()),
829 rg_mode,
830 pre_amp_db,
831 move || {
832 finish_tx.send(PlayerCommand::DecodeFinished).ok();
833 },
834 )?;
835
836 let (engine, rate_watch) = self.create_engine_for(&info, consumer)?;
837 engine.start()?;
838
839 self.shared_state.set_playback_state(PlaybackState::Playing);
840
841 self.active_playback = Some(ActivePlayback {
842 engine,
843 decode_handle,
844 stream: Some(live),
845 _rate_watch: rate_watch,
846 });
847
848 Ok(())
849 }
850
851 pub fn seek(&mut self, position_ms: u64) {
858 let Some(info) = self.shared_state.track_info() else {
859 return;
860 };
861 let seekable = self.shared_state.seekable_ms();
863 if seekable == 0 {
864 log::debug!("seek declined: {:?} is not seekable yet", info.id);
868 return;
869 }
870 let ceiling = seekable.min(
871 self.shared_state
872 .duration_ms()
873 .saturating_sub(SEEK_END_GUARD_MS),
874 );
875 let clamped = position_ms.min(ceiling);
876
877 if let Err(e) = self.restart_current(&info, clamped) {
878 log::error!("seek failed: {}", e);
879 }
880 }
881
882 fn restart_current(&mut self, info: &TrackInfo, position_ms: u64) -> Result<(), PlayerError> {
890 let was_paused = self.shared_state.playback_state() == PlaybackState::Paused;
891
892 match self.shared_state.item_playback_source(info.id) {
893 Some(PlaybackSource::Streaming {
894 path,
895 bytes_written,
896 total,
897 }) => {
898 let known = buffer::StreamInfo {
902 codec: info.codec.clone(),
903 sample_rate: info.sample_rate,
904 channels: info.channels,
905 bit_depth: info.bit_depth,
906 bitrate_kbps: info.bitrate_kbps,
907 duration_ms: info.duration_ms,
908 };
909 let source = StreamSource {
910 path,
911 bytes_written,
912 total,
913 mode: self.stream_mode,
914 };
915 self.start_streaming_playback(info.id, source, position_ms, known)?;
916 }
917 Some(PlaybackSource::Ready(path)) => {
918 self.start_playback(info.id, &path, position_ms)?;
919 }
920 None => return Ok(()),
921 }
922
923 if was_paused {
924 self.pause_now();
925 }
926 Ok(())
927 }
928
929 pub fn next_track(&mut self) {
931 match self.shared_state.advance_cursor_loadable() {
932 Some(id) => self.play(id),
933 None => {
934 log::info!("no more tracks in playlist");
935 self.stop_playback_and_clear_state();
936 }
937 }
938 }
939
940 pub fn prev_track(&mut self) {
942 match self.shared_state.retreat_cursor() {
943 Some((id, _)) => self.play(id),
944 None => {
945 if let Some(info) = self.shared_state.track_info()
947 && let Err(e) = self.restart_current(&info, 0)
948 {
949 log::error!("restart failed: {}", e);
950 }
951 }
952 }
953 }
954
955 pub fn pause(&mut self) {
960 let Some(ref playback) = self.active_playback else {
961 return;
962 };
963 if crate::config::Config::load_or_default()
964 .playback
965 .fade_on_pause
966 {
967 playback.engine.fade_out();
968 self.shared_state.set_playback_state(PlaybackState::Paused);
969 } else {
970 self.pause_now();
971 }
972 }
973
974 fn pause_now(&mut self) {
977 if let Some(ref playback) = self.active_playback {
978 if let Err(e) = playback.engine.stop() {
979 log::error!("pause failed: {}", e);
980 return;
981 }
982 self.shared_state.set_playback_state(PlaybackState::Paused);
983 }
984 }
985
986 pub fn resume(&mut self) {
988 if let Some(ref playback) = self.active_playback {
989 let engine = &playback.engine;
990 let resumed = if engine.is_running() || engine.is_silent() {
991 engine.fade_in()
992 } else {
993 engine.start()
994 };
995 if let Err(e) = resumed {
996 log::error!("resume failed: {}", e);
997 return;
998 }
999 self.shared_state.set_playback_state(PlaybackState::Playing);
1000 self.wake_analyzer();
1001 }
1002 }
1003
1004 fn wake_analyzer(&self) {
1011 self.viz_snapshot.wake();
1012 }
1013
1014 pub fn stop(&mut self) {
1016 self.shared_state.clear_playlist();
1017 self.stop_playback_and_clear_state();
1018 }
1019
1020 fn stop_engine(&mut self) {
1026 let Some(playback) = self.active_playback.take() else {
1027 return;
1028 };
1029 let ActivePlayback {
1030 engine,
1031 mut decode_handle,
1032 stream,
1033 _rate_watch,
1034 } = playback;
1035
1036 let _ = engine.stop();
1037 decode_handle.signal_stop();
1040 if let Some(stream) = stream {
1041 stream.abandon();
1042 }
1043 decode_handle.stop();
1044 drop(engine);
1045 }
1046
1047 fn stop_playback_and_clear_state(&mut self) {
1049 self.finish_play();
1050 self.stop_engine();
1051 self.timeline.reset();
1052 self.shared_state.set_playback_state(PlaybackState::Stopped);
1053 self.shared_state.set_position_ms(0);
1054 self.shared_state.set_track_info(None);
1055 }
1056
1057 pub fn remove_from_playlist(&mut self, id: QueueItemId) {
1064 let was_cursor = self.shared_state.is_cursor(id);
1065 let resume_after = was_cursor
1066 .then(|| self.shared_state.item_before(id))
1067 .flatten();
1068 self.shared_state.remove_item(id);
1069 if was_cursor {
1070 self.shared_state.set_cursor(resume_after);
1071 self.next_track();
1072 }
1073 }
1074
1075 pub fn track_ready(&mut self, id: QueueItemId) {
1078 self.shared_state.update_item_state(id, ItemState::Ready);
1080
1081 if !self.shared_state.is_cursor(id) {
1082 return;
1083 }
1084
1085 let is_playing = self.shared_state.playback_state() == PlaybackState::Playing;
1086 let current_track_id = self.shared_state.track_info().map(|t| t.id);
1087
1088 if is_playing && current_track_id == Some(id) {
1089 log::info!(
1092 "track_ready: download complete while streaming {:?}, refreshing metadata",
1093 id
1094 );
1095 self.refresh_track_metadata(id);
1096 return;
1097 }
1098
1099 if !is_playing && let Some(path) = self.shared_state.item_path_if_ready(id) {
1101 log::info!("track_ready: starting playback for {:?}", id);
1102 if let Err(e) = self.start_playback(id, &path, 0) {
1103 log::error!("track_ready playback failed: {}", e);
1104 }
1105 }
1106 }
1107
1108 pub fn track_stream_ready(&mut self, id: QueueItemId) {
1111 if !self.shared_state.is_cursor(id) {
1112 return;
1113 }
1114
1115 let is_playing = self.shared_state.playback_state() == PlaybackState::Playing;
1116 if is_playing {
1117 return; }
1119
1120 match self.shared_state.item_playback_source(id) {
1121 Some(PlaybackSource::Streaming {
1122 path,
1123 bytes_written,
1124 total,
1125 }) => {
1126 log::info!("track_stream_ready: probing partial file for {:?}", id);
1127 self.probe_stream_for_playback(id, &path, bytes_written, total);
1128 }
1129 Some(PlaybackSource::Ready(path)) => {
1130 log::info!(
1132 "track_stream_ready: track already ready, starting normal playback for {:?}",
1133 id
1134 );
1135 if let Err(e) = self.start_playback(id, &path, 0) {
1136 log::error!("track_stream_ready playback failed: {}", e);
1137 }
1138 }
1139 None => {} }
1141 }
1142
1143 fn refresh_track_metadata(&mut self, id: QueueItemId) {
1147 use crate::index::metadata;
1148
1149 let path = match self.shared_state.item_path_if_ready(id) {
1150 Some(p) => p,
1151 None => return,
1152 };
1153
1154 match metadata::read_metadata(&path) {
1155 Ok(meta) => {
1156 self.shared_state.update_item_metadata(
1157 id,
1158 meta.title,
1159 meta.artist,
1160 meta.album_artist.unwrap_or_default(),
1161 meta.album,
1162 meta.duration_ms.map(|d| d as u64),
1163 );
1164
1165 if let Some(current) = self.shared_state.track_info()
1174 && current.id == id
1175 {
1176 let probed = buffer::probe_file(&path).ok();
1177 let duration_ms = probed
1178 .as_ref()
1179 .map(|s| s.duration_ms)
1180 .filter(|d| *d > current.duration_ms)
1181 .unwrap_or(current.duration_ms);
1182 if duration_ms != current.duration_ms {
1183 log::info!(
1184 "track_ready: duration corrected {}ms → {}ms",
1185 current.duration_ms,
1186 duration_ms
1187 );
1188 }
1189 self.shared_state.set_track_info(Some(TrackInfo {
1190 duration_ms,
1191 path: path.clone(),
1192 ..current
1193 }));
1194 }
1195
1196 self.shared_state.signal_metadata_refresh();
1198 log::info!("track_ready: metadata refreshed for {:?}", id);
1199 }
1200 Err(e) => {
1201 log::warn!("track_ready: metadata refresh failed for {:?}: {}", id, e);
1202 }
1203 }
1204 }
1205
1206 fn on_track_changed(&mut self, id: QueueItemId) {
1216 if self.in_flight.as_ref().is_some_and(|f| f.item == id) {
1217 return;
1218 }
1219 self.finish_play();
1220 let track_id = self.shared_state.item_db_id(id);
1221 self.in_flight = Some(InFlight::new(id, track_id));
1222 if let (Some(track_id), Some(recorder)) = (track_id, self.history.as_ref()) {
1223 recorder.record(PlayEvent::Started { track_id });
1224 }
1225 }
1226
1227 fn finish_play(&mut self) -> Option<PlayEvent> {
1230 let flight = self.in_flight.take()?;
1231 let event = PlayEvent::Finished {
1232 track_id: flight.track_id()?,
1233 listened_ms: flight.listened_ms(),
1234 };
1235 if let Some(recorder) = self.history.as_ref() {
1236 recorder.record(event);
1237 }
1238 Some(event)
1239 }
1240
1241 pub fn update_playback_state(&mut self) {
1242 let Some(playback) = self.active_playback.as_ref() else {
1243 return;
1244 };
1245
1246 if self.shared_state.playback_state() == PlaybackState::Paused
1247 && playback.engine.is_running()
1248 && playback.engine.is_silent()
1249 && let Err(e) = playback.engine.stop()
1250 {
1251 log::error!("stopping after fade failed: {}", e);
1252 }
1253
1254 if let Some((id, path, info, position_ms)) = self.timeline.current_playback() {
1255 self.shared_state.set_position_ms(position_ms);
1256
1257 self.on_track_changed(id);
1260 if let Some(f) = self.in_flight.as_mut() {
1261 f.advance(position_ms);
1262 }
1263
1264 let current_id = self.shared_state.track_info().map(|t| t.id);
1267 if current_id != Some(id) {
1268 log::info!("timeline: now playing {:?}", id);
1269 self.shared_state.set_track_info(Some(TrackInfo {
1270 id,
1271 path,
1272 codec: info.codec,
1273 sample_rate: info.sample_rate,
1274 bit_depth: info.bit_depth,
1275 bitrate_kbps: info.bitrate_kbps,
1276 channels: info.channels,
1277 duration_ms: info.duration_ms,
1278 }));
1279 self.shared_state.set_cursor(Some(id));
1280 }
1281 }
1282 }
1283
1284 pub fn track_failed(&mut self, id: QueueItemId) {
1291 if !self.shared_state.is_cursor(id) {
1292 return;
1293 }
1294 if self.shared_state.playback_state() != PlaybackState::Stopped {
1298 return;
1299 }
1300 log::info!("track {:?} cannot load, moving on", id);
1301 self.next_track();
1302 }
1303
1304 fn on_decode_finished(&mut self) {
1311 log::info!("decode finished, checking for next track");
1312 match self.shared_state.advance_cursor_loadable() {
1313 Some(id) => self.play(id),
1314 None => {
1315 log::info!("no more tracks — stopping");
1316 self.stop_playback_and_clear_state();
1317 }
1318 }
1319 }
1320
1321 fn snapshot_for_undo(
1325 &self,
1326 ids: &[QueueItemId],
1327 ) -> Vec<(Box<state::PlaylistItem>, Option<QueueItemId>)> {
1328 self.shared_state
1329 .items_before(ids)
1330 .into_iter()
1331 .filter_map(|(id, after)| Some((Box::new(self.shared_state.get_item(id)?), after)))
1332 .collect()
1333 }
1334
1335 fn push_undo(&mut self, entry: UndoEntry) {
1337 if let Some(ref mut batch) = self.batch_buffer {
1338 batch.push(entry);
1339 } else {
1340 self.undo_stack.push(entry);
1341 }
1342 }
1343
1344 pub fn process_command(&mut self, cmd: PlayerCommand) {
1346 match cmd {
1347 PlayerCommand::Play(id) => self.play(id),
1348 PlayerCommand::Pause => self.pause(),
1349 PlayerCommand::Resume => self.resume(),
1350 PlayerCommand::Stop => self.stop(),
1351 PlayerCommand::Seek(pos) => self.seek(pos),
1352 PlayerCommand::NextTrack => {
1353 let now = std::time::Instant::now();
1355 if now.duration_since(self.last_skip).as_millis() >= 150 {
1356 self.last_skip = now;
1357 self.next_track();
1358 }
1359 }
1360 PlayerCommand::PrevTrack => {
1361 let now = std::time::Instant::now();
1362 if now.duration_since(self.last_skip).as_millis() >= 150 {
1363 self.last_skip = now;
1364 self.prev_track();
1365 }
1366 }
1367 PlayerCommand::AddToPlaylist(items) => {
1368 let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
1369 self.shared_state.add_items(items);
1370 self.push_undo(UndoEntry::Added { ids });
1371 }
1372 PlayerCommand::UpdatePaths(updates) => {
1373 self.shared_state.update_paths(&updates);
1374 if let Some(info) = self.shared_state.track_info()
1375 && let Some((_, new_path)) = updates.iter().find(|(id, _)| *id == info.id)
1376 {
1377 self.shared_state.set_track_info(Some(TrackInfo {
1378 path: new_path.clone(),
1379 ..info
1380 }));
1381 }
1382 }
1383 PlayerCommand::InsertInPlaylist { items, after } => {
1384 let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
1385 self.shared_state.insert_items_after(items, after);
1386 self.push_undo(UndoEntry::Inserted { ids });
1387 }
1388 PlayerCommand::ClearPlaylist => {
1389 self.stop_playback_and_clear_state();
1393 let (items, cursor) = self.shared_state.snapshot_playlist();
1394 self.shared_state.clear_playlist();
1395 self.push_undo(UndoEntry::Replaced { items, cursor });
1396 }
1397 PlayerCommand::ReplacePlaylist { items, start } => {
1398 self.stop_playback_and_clear_state();
1402 let (old_items, cursor) = self.shared_state.snapshot_playlist();
1403 self.shared_state.clear_playlist();
1404 self.push_undo(UndoEntry::Replaced {
1405 items: old_items,
1406 cursor,
1407 });
1408
1409 if items.is_empty() {
1410 return;
1411 }
1412 let start_id = items.get(start).unwrap_or(&items[0]).id;
1413 self.shared_state.add_items(items);
1414 self.play(start_id);
1415 }
1416 PlayerCommand::RemoveFromPlaylist(id) => {
1417 let item = self.shared_state.get_item(id);
1418 let after = self.shared_state.item_before(id);
1419 self.remove_from_playlist(id);
1420 if let Some(item) = item {
1421 self.push_undo(UndoEntry::Removed {
1422 items: vec![(Box::new(item), after)],
1423 });
1424 }
1425 }
1426 PlayerCommand::RemoveFromPlaylistBatch(ids) => {
1427 let items_with_pos = self.snapshot_for_undo(&ids);
1431 let resume_after = match self.shared_state.cursor() {
1432 Some(cursor) if ids.contains(&cursor) => {
1433 Some(self.shared_state.surviving_item_before(cursor, &ids))
1434 }
1435 _ => None,
1436 };
1437
1438 self.shared_state.remove_items(&ids);
1439
1440 if let Some(resume_after) = resume_after {
1441 self.shared_state.set_cursor(resume_after);
1442 self.next_track();
1443 }
1444
1445 if !items_with_pos.is_empty() {
1446 self.push_undo(UndoEntry::Removed {
1447 items: items_with_pos,
1448 });
1449 }
1450 }
1451 PlayerCommand::MoveInPlaylist { id, target, after } => {
1452 let was_after = self.shared_state.item_before(id);
1453 self.shared_state.move_item(id, target, after);
1454 self.push_undo(UndoEntry::Moved { id, was_after });
1455 }
1456 PlayerCommand::MoveItemsInPlaylist { ids, target, after } => {
1457 let entries = self.shared_state.items_before(&ids);
1458 self.shared_state.move_items(&ids, target, after);
1459 self.push_undo(UndoEntry::MovedBatch { entries });
1460 }
1461 PlayerCommand::ReorderPlaylist(order) => {
1462 let entries = self.shared_state.items_before(&order);
1466 self.shared_state.reorder_to(&order);
1467 self.push_undo(UndoEntry::MovedBatch { entries });
1468 }
1469 PlayerCommand::TrackReady(id) => self.track_ready(id),
1470 PlayerCommand::DecodeFinished => self.on_decode_finished(),
1471 PlayerCommand::TrackStreamReady(id) => self.track_stream_ready(id),
1472 PlayerCommand::StreamProbed { id, info, mode } => self.stream_probed(id, *info, mode),
1473 PlayerCommand::TrackFailed(id) => self.track_failed(id),
1474 PlayerCommand::Undo => self.execute_undo(),
1475 PlayerCommand::Redo => self.execute_redo(),
1476 PlayerCommand::BeginUndoBatch => {
1477 self.batch_buffer = Some(Vec::new());
1478 }
1479 PlayerCommand::EndUndoBatch => {
1480 if let Some(entries) = self.batch_buffer.take() {
1481 if entries.len() == 1 {
1482 self.undo_stack.push(entries.into_iter().next().unwrap());
1484 } else if !entries.is_empty() {
1485 self.undo_stack.push(UndoEntry::Batch(entries));
1486 }
1487 }
1488 }
1489 PlayerCommand::SetOutputDevice(name) => self.set_output_device(name),
1490 PlayerCommand::ClearOutputDevice => self.clear_output_device(),
1491 }
1492 }
1493
1494 fn apply_entry(&mut self, entry: UndoEntry) -> Option<UndoEntry> {
1496 match entry {
1497 UndoEntry::Added { ids } => {
1498 let items_with_pos = self.snapshot_for_undo(&ids);
1500 self.shared_state.remove_items(&ids);
1501 Some(UndoEntry::Removed {
1502 items: items_with_pos,
1503 })
1504 }
1505 UndoEntry::Removed { items } => {
1506 let mut ids = Vec::with_capacity(items.len());
1508 for (item, after) in items {
1509 ids.push(item.id);
1510 self.shared_state.insert_item_at(*item, after);
1511 }
1512 Some(UndoEntry::Added { ids })
1513 }
1514 UndoEntry::Inserted { ids } => {
1515 let items_with_pos = self.snapshot_for_undo(&ids);
1517 self.shared_state.remove_items(&ids);
1518 Some(UndoEntry::Removed {
1519 items: items_with_pos,
1520 })
1521 }
1522 UndoEntry::Moved { id, was_after } => {
1523 let current_after = self.shared_state.item_before(id);
1524 self.shared_state.move_item_to(id, was_after);
1525 Some(UndoEntry::Moved {
1526 id,
1527 was_after: current_after,
1528 })
1529 }
1530 UndoEntry::MovedBatch { entries } => {
1531 let ids: Vec<QueueItemId> = entries.iter().map(|(id, _)| *id).collect();
1532 let current_positions = self.shared_state.items_before(&ids);
1533 self.shared_state.move_items_to(&entries);
1534 Some(UndoEntry::MovedBatch {
1535 entries: current_positions,
1536 })
1537 }
1538 UndoEntry::Replaced { items, cursor } => {
1539 let (current_items, current_cursor) = self.shared_state.snapshot_playlist();
1540 self.shared_state.restore_playlist(items, cursor);
1541 Some(UndoEntry::Replaced {
1542 items: current_items,
1543 cursor: current_cursor,
1544 })
1545 }
1546 UndoEntry::Batch(entries) => {
1547 let mut inverses = Vec::with_capacity(entries.len());
1549 for entry in entries.into_iter().rev() {
1550 if let Some(inverse) = self.apply_entry(entry) {
1551 inverses.push(inverse);
1552 }
1553 }
1554 inverses.reverse();
1555 Some(UndoEntry::Batch(inverses))
1556 }
1557 }
1558 }
1559
1560 fn reconcile_playback(&mut self) {
1574 let Some(playing) = self.shared_state.track_info().map(|t| t.id) else {
1575 return;
1576 };
1577 if self.shared_state.get_item(playing).is_some() {
1578 return;
1579 }
1580 let resume = (self.shared_state.playback_state() == PlaybackState::Playing)
1585 .then(|| self.shared_state.cursor())
1586 .flatten();
1587 self.stop_playback_and_clear_state();
1588 if let Some(id) = resume {
1589 self.play(id);
1590 }
1591 }
1592
1593 fn execute_undo(&mut self) {
1595 let Some(entry) = self.undo_stack.pop_undo() else {
1596 return;
1597 };
1598 if let Some(inverse) = self.apply_entry(entry) {
1599 self.undo_stack.push_redo(inverse);
1600 }
1601 self.reconcile_playback();
1602 }
1603
1604 fn execute_redo(&mut self) {
1606 let Some(entry) = self.undo_stack.pop_redo() else {
1607 return;
1608 };
1609 if let Some(inverse) = self.apply_entry(entry) {
1610 self.undo_stack.push_undo_keep_redo(inverse);
1611 }
1612 self.reconcile_playback();
1613 }
1614
1615 pub fn run(&mut self) {
1617 use std::time::Duration;
1618
1619 let rx = self.commands.rx.clone();
1620 loop {
1621 match rx.recv_timeout(Duration::from_millis(50)) {
1623 Ok(cmd) => self.process_command(cmd),
1624 Err(crossbeam_channel::RecvTimeoutError::Timeout) => {}
1625 Err(crossbeam_channel::RecvTimeoutError::Disconnected) => break,
1626 }
1627 self.update_playback_state();
1628 }
1629 self.stop();
1630 }
1631
1632 pub fn spawn() -> (
1635 Arc<SharedPlayerState>,
1636 Arc<PlaybackTimeline>,
1637 Arc<VizSnapshot>,
1638 crossbeam_channel::Sender<PlayerCommand>,
1639 ) {
1640 let mut player = Self::new();
1641 player.history = PlayRecorder::spawn();
1642 let state = player.shared_state();
1643 let timeline = player.timeline();
1644 let viz_snapshot = player.viz_snapshot();
1645 let tx = player.command_sender();
1646
1647 thread::Builder::new()
1648 .name("koan-player".into())
1649 .spawn(move || player.run())
1650 .expect("failed to spawn player thread");
1651
1652 (state, timeline, viz_snapshot, tx)
1653 }
1654}
1655
1656#[cfg(test)]
1657mod tests {
1658 use super::*;
1659 use state::PlaylistItem;
1660 use std::path::PathBuf;
1661 use std::sync::atomic::AtomicU64;
1662
1663 fn make_item(title: &str) -> PlaylistItem {
1664 PlaylistItem {
1665 playlist_entry_id: None,
1666 id: QueueItemId::new(),
1667 db_id: None,
1668 path: PathBuf::from(format!("/music/{title}.flac")),
1669 title: title.to_string(),
1670 artist: String::new(),
1671 album_artist: String::new(),
1672 album: String::new(),
1673 year: None,
1674 codec: None,
1675 track_number: None,
1676 disc: None,
1677 duration_ms: None,
1678 state: ItemState::Ready,
1679 }
1680 }
1681
1682 fn playlist_ids(player: &Player) -> Vec<QueueItemId> {
1683 let (items, _) = player.shared_state.snapshot_playlist();
1684 items.iter().map(|i| i.id).collect()
1685 }
1686
1687 fn playlist_titles(player: &Player) -> Vec<String> {
1688 let (items, _) = player.shared_state.snapshot_playlist();
1689 items.iter().map(|i| i.title.clone()).collect()
1690 }
1691
1692 fn pending_item(title: &str) -> PlaylistItem {
1693 PlaylistItem {
1694 playlist_entry_id: None,
1695 state: ItemState::Pending,
1696 ..make_item(title)
1697 }
1698 }
1699
1700 fn pretend_playing(player: &mut Player, id: QueueItemId) {
1704 let item = player
1705 .shared_state
1706 .get_item(id)
1707 .expect("item is in the queue");
1708 player.shared_state.set_track_info(Some(TrackInfo {
1709 id,
1710 path: item.path,
1711 codec: String::new(),
1712 sample_rate: 44_100,
1713 bit_depth: None,
1714 bitrate_kbps: None,
1715 channels: 2,
1716 duration_ms: 1_000,
1717 }));
1718 player
1719 .shared_state
1720 .set_playback_state(PlaybackState::Playing);
1721 }
1722
1723 fn playing_id(player: &Player) -> Option<QueueItemId> {
1724 player.shared_state.track_info().map(|t| t.id)
1725 }
1726
1727 fn seed(player: &mut Player, n: usize) -> Vec<QueueItemId> {
1729 let items: Vec<_> = (0..n).map(|i| make_item(&format!("t{i}"))).collect();
1730 let ids = items.iter().map(|i| i.id).collect();
1731 player.process_command(PlayerCommand::AddToPlaylist(items));
1732 ids
1733 }
1734
1735 fn listen(player: &mut Player, from_ms: u64, to_ms: u64) {
1739 let mut at = from_ms;
1740 if let Some(f) = player.in_flight.as_mut() {
1741 f.advance(at); }
1743 while at < to_ms {
1744 at = (at + 50).min(to_ms);
1745 if let Some(f) = player.in_flight.as_mut() {
1746 f.advance(at);
1747 }
1748 }
1749 }
1750
1751 fn start(player: &mut Player, track_id: i64) -> QueueItemId {
1752 let id = QueueItemId::new();
1753 player.on_track_changed(id);
1754 player
1756 .in_flight
1757 .as_mut()
1758 .unwrap()
1759 .track_id_for_test(track_id);
1760 id
1761 }
1762
1763 #[test]
1764 fn a_gapless_transition_closes_the_outgoing_track_and_opens_the_next() {
1765 let mut player = Player::new();
1766 start(&mut player, 11);
1767 listen(&mut player, 0, 200_000);
1768
1769 let b = QueueItemId::new();
1770 player.on_track_changed(b);
1771 let f = player
1772 .in_flight
1773 .as_ref()
1774 .expect("the next track is counting");
1775 assert_eq!(f.item, b);
1776 assert_eq!(f.listened_ms(), 0, "and starts from nothing");
1777 }
1778
1779 #[test]
1780 fn a_track_skipped_seconds_in_is_still_history() {
1781 let mut player = Player::new();
1782 start(&mut player, 7);
1783 listen(&mut player, 0, 2_000);
1784
1785 let event = player
1786 .finish_play()
1787 .expect("putting something on is a thing you did, however briefly");
1788 assert!(matches!(
1789 event,
1790 history::PlayEvent::Finished {
1791 track_id: 7,
1792 listened_ms: 2_000
1793 }
1794 ));
1795 }
1796
1797 #[test]
1798 fn a_track_is_closed_out_once() {
1799 let mut player = Player::new();
1800 start(&mut player, 7);
1801 listen(&mut player, 0, 200_000);
1802
1803 assert!(player.finish_play().is_some());
1804 assert!(player.finish_play().is_none());
1805 }
1806
1807 #[test]
1808 fn seeking_around_a_track_does_not_enter_it_twice() {
1809 let mut player = Player::new();
1810 let id = start(&mut player, 7);
1811 listen(&mut player, 0, 120_000);
1812
1813 player.on_track_changed(id);
1815 assert_eq!(
1816 player.in_flight.as_ref().unwrap().listened_ms(),
1817 120_000,
1818 "the seek kept the count rather than restarting it"
1819 );
1820 listen(&mut player, 30_000, 40_000);
1821
1822 let Some(history::PlayEvent::Finished { listened_ms, .. }) = player.finish_play() else {
1823 panic!("still one play");
1824 };
1825 assert_eq!(listened_ms, 130_000);
1826 assert!(player.finish_play().is_none());
1827 }
1828
1829 #[test]
1830 fn a_track_that_is_not_in_the_library_is_not_recorded() {
1831 let mut player = Player::new();
1832 let id = QueueItemId::new();
1833 player.on_track_changed(id);
1834 listen(&mut player, 0, 200_000);
1835 assert!(player.finish_play().is_none());
1836 }
1837
1838 #[test]
1839 fn stopping_closes_out_what_was_heard() {
1840 let mut player = Player::new();
1841 start(&mut player, 7);
1842 listen(&mut player, 0, 150_000);
1843
1844 player.stop_playback_and_clear_state();
1845 assert!(player.in_flight.is_none(), "the stop consumed it");
1846 }
1847
1848 #[test]
1849 fn removing_the_playing_track_resumes_at_its_successor() {
1850 let mut player = Player::new();
1851 let ids = seed(&mut player, 5);
1852 player.shared_state.set_cursor(Some(ids[2]));
1853
1854 player.process_command(PlayerCommand::RemoveFromPlaylist(ids[2]));
1855
1856 assert_eq!(
1857 player.shared_state.cursor(),
1858 Some(ids[3]),
1859 "playback must continue at the next track, not restart the queue"
1860 );
1861 assert_eq!(player.playback_starts, 1);
1862 }
1863
1864 #[test]
1865 fn removing_the_first_playing_track_resumes_at_the_new_first() {
1866 let mut player = Player::new();
1867 let ids = seed(&mut player, 3);
1868 player.shared_state.set_cursor(Some(ids[0]));
1869
1870 player.process_command(PlayerCommand::RemoveFromPlaylist(ids[0]));
1871
1872 assert_eq!(player.shared_state.cursor(), Some(ids[1]));
1873 }
1874
1875 #[test]
1876 fn next_track_parks_on_a_track_that_has_not_downloaded_yet() {
1877 let mut player = Player::new();
1878 let playing = make_item("playing");
1879 let waiting = pending_item("waiting");
1880 let later = make_item("later");
1881 let (playing_id, waiting_id) = (playing.id, waiting.id);
1882 player.process_command(PlayerCommand::AddToPlaylist(vec![playing, waiting, later]));
1883 player.shared_state.set_cursor(Some(playing_id));
1884
1885 player.process_command(PlayerCommand::DecodeFinished);
1886
1887 assert_eq!(
1888 player.shared_state.cursor(),
1889 Some(waiting_id),
1890 "the cursor parks on the track being fetched"
1891 );
1892 assert_eq!(
1893 player.playback_starts, 0,
1894 "nothing to play until its bytes land"
1895 );
1896
1897 player
1900 .shared_state
1901 .update_item_state(waiting_id, ItemState::Ready);
1902 player.process_command(PlayerCommand::TrackReady(waiting_id));
1903
1904 assert_eq!(player.playback_starts, 1);
1905 assert_eq!(player.shared_state.cursor(), Some(waiting_id));
1906 }
1907
1908 #[test]
1909 fn a_download_that_cannot_land_moves_the_cursor_on() {
1910 let mut player = Player::new();
1911 let waiting = pending_item("waiting");
1912 let later = make_item("later");
1913 let (waiting_id, later_id) = (waiting.id, later.id);
1914 player.process_command(PlayerCommand::AddToPlaylist(vec![waiting, later]));
1915
1916 player.process_command(PlayerCommand::Play(waiting_id));
1917 assert_eq!(player.playback_starts, 0, "nothing to play yet");
1918
1919 player
1921 .shared_state
1922 .update_item_state(waiting_id, ItemState::Failed("remote unavailable".into()));
1923 player.process_command(PlayerCommand::TrackFailed(waiting_id));
1924
1925 assert_eq!(
1926 player.shared_state.cursor(),
1927 Some(later_id),
1928 "the queue moves past a track that can never load"
1929 );
1930 assert_eq!(player.playback_starts, 1);
1931 }
1932
1933 #[test]
1934 fn a_queue_that_can_never_load_stops_rather_than_waiting() {
1935 let mut player = Player::new();
1936 let first = pending_item("first");
1937 let second = pending_item("second");
1938 let (first_id, second_id) = (first.id, second.id);
1939 player.process_command(PlayerCommand::AddToPlaylist(vec![first, second]));
1940
1941 player.process_command(PlayerCommand::Play(first_id));
1942 for id in [first_id, second_id] {
1943 player
1944 .shared_state
1945 .update_item_state(id, ItemState::Failed("remote unavailable".into()));
1946 player.process_command(PlayerCommand::TrackFailed(id));
1947 }
1948
1949 assert_eq!(player.playback_starts, 0);
1950 assert_eq!(
1951 player.shared_state.playback_state(),
1952 PlaybackState::Stopped,
1953 "a stop the UI can see, not an indefinite wait for TrackReady"
1954 );
1955 }
1956
1957 #[test]
1958 fn a_failure_elsewhere_in_the_queue_leaves_the_cursor_alone() {
1959 let mut player = Player::new();
1960 let waiting = pending_item("waiting");
1961 let other = pending_item("other");
1962 let (waiting_id, other_id) = (waiting.id, other.id);
1963 player.process_command(PlayerCommand::AddToPlaylist(vec![waiting, other]));
1964 player.process_command(PlayerCommand::Play(waiting_id));
1965
1966 player
1967 .shared_state
1968 .update_item_state(other_id, ItemState::Failed("remote unavailable".into()));
1969 player.process_command(PlayerCommand::TrackFailed(other_id));
1970
1971 assert_eq!(
1972 player.shared_state.cursor(),
1973 Some(waiting_id),
1974 "a track still downloading keeps the cursor"
1975 );
1976 }
1977
1978 #[test]
1979 fn batch_delete_containing_the_cursor_restarts_the_engine_once() {
1980 let mut player = Player::new();
1981 let ids = seed(&mut player, 5);
1982 player.shared_state.set_cursor(Some(ids[2]));
1983
1984 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![
1985 ids[1], ids[2], ids[3],
1986 ]));
1987
1988 assert_eq!(playlist_titles(&player), vec!["t0", "t4"]);
1989 assert_eq!(player.shared_state.cursor(), Some(ids[4]));
1990 assert_eq!(
1991 player.playback_starts, 1,
1992 "one resume for the whole selection, not one per deleted track"
1993 );
1994 }
1995
1996 #[test]
1997 fn batch_delete_below_the_cursor_leaves_playback_alone() {
1998 let mut player = Player::new();
1999 let ids = seed(&mut player, 4);
2000 player.shared_state.set_cursor(Some(ids[0]));
2001
2002 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![ids[2], ids[3]]));
2003
2004 assert_eq!(player.shared_state.cursor(), Some(ids[0]));
2005 assert_eq!(player.playback_starts, 0);
2006 }
2007
2008 #[test]
2009 fn undo_of_a_batch_delete_restores_the_original_order() {
2010 let mut player = Player::new();
2014 let items = vec![
2015 make_item("A"),
2016 make_item("B"),
2017 make_item("C"),
2018 make_item("D"),
2019 ];
2020 let (b_id, c_id) = (items[1].id, items[2].id);
2021 player.process_command(PlayerCommand::AddToPlaylist(items));
2022
2023 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![c_id, b_id]));
2024 assert_eq!(playlist_titles(&player), vec!["A", "D"]);
2025
2026 player.process_command(PlayerCommand::Undo);
2027 assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
2028 }
2029
2030 #[test]
2033 fn undo_add_removes_items() {
2034 let mut player = Player::new();
2035 let items = vec![make_item("A"), make_item("B")];
2036 let ids: Vec<_> = items.iter().map(|i| i.id).collect();
2037
2038 player.process_command(PlayerCommand::AddToPlaylist(items));
2039 assert_eq!(playlist_ids(&player), ids);
2040 assert!(player.undo_stack().can_undo());
2041
2042 player.process_command(PlayerCommand::Undo);
2043 assert!(playlist_ids(&player).is_empty());
2044 assert!(player.undo_stack().can_redo());
2045 }
2046
2047 #[test]
2048 fn redo_add_restores_items() {
2049 let mut player = Player::new();
2050 let items = vec![make_item("A"), make_item("B")];
2051
2052 player.process_command(PlayerCommand::AddToPlaylist(items));
2053 player.process_command(PlayerCommand::Undo);
2054 assert!(playlist_ids(&player).is_empty());
2055
2056 player.process_command(PlayerCommand::Redo);
2057 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2058 }
2059
2060 #[test]
2063 fn undo_remove_restores_item_at_position() {
2064 let mut player = Player::new();
2065 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2066 let b_id = items[1].id;
2067
2068 player.process_command(PlayerCommand::AddToPlaylist(items));
2069 player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
2070 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2071
2072 player.process_command(PlayerCommand::Undo);
2073 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2074 }
2075
2076 #[test]
2077 fn undo_remove_first_item() {
2078 let mut player = Player::new();
2079 let items = vec![make_item("A"), make_item("B")];
2080 let a_id = items[0].id;
2081
2082 player.process_command(PlayerCommand::AddToPlaylist(items));
2083 player.process_command(PlayerCommand::RemoveFromPlaylist(a_id));
2084 assert_eq!(playlist_titles(&player), vec!["B"]);
2085
2086 player.process_command(PlayerCommand::Undo);
2087 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2088 }
2089
2090 #[test]
2091 fn undo_batch_remove_restores_all() {
2092 let mut player = Player::new();
2093 let items = vec![
2094 make_item("A"),
2095 make_item("B"),
2096 make_item("C"),
2097 make_item("D"),
2098 ];
2099 let b_id = items[1].id;
2100 let c_id = items[2].id;
2101
2102 player.process_command(PlayerCommand::AddToPlaylist(items));
2103 let version_before = player.shared_state.playlist_version();
2104 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![b_id, c_id]));
2105 assert_eq!(playlist_titles(&player), vec!["A", "D"]);
2106 assert_eq!(
2109 player.shared_state.playlist_version(),
2110 version_before + 1,
2111 "batch removal must bump the playlist version exactly once"
2112 );
2113
2114 player.process_command(PlayerCommand::Undo);
2116 assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
2117 }
2118
2119 #[test]
2120 fn redo_batch_remove() {
2121 let mut player = Player::new();
2122 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2123 let a_id = items[0].id;
2124 let b_id = items[1].id;
2125
2126 player.process_command(PlayerCommand::AddToPlaylist(items));
2127 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![a_id, b_id]));
2128 player.process_command(PlayerCommand::Undo);
2129 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2130
2131 player.process_command(PlayerCommand::Redo);
2132 assert_eq!(playlist_titles(&player), vec!["C"]);
2133 }
2134
2135 #[test]
2136 fn redo_remove() {
2137 let mut player = Player::new();
2138 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2139 let b_id = items[1].id;
2140
2141 player.process_command(PlayerCommand::AddToPlaylist(items));
2142 player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
2143 player.process_command(PlayerCommand::Undo);
2144 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2145
2146 player.process_command(PlayerCommand::Redo);
2147 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2148 }
2149
2150 #[test]
2153 fn undo_insert_removes_inserted_items() {
2154 let mut player = Player::new();
2155 let items = vec![make_item("A"), make_item("C")];
2156 let a_id = items[0].id;
2157
2158 player.process_command(PlayerCommand::AddToPlaylist(items));
2159
2160 let inserted = vec![make_item("B")];
2161 player.process_command(PlayerCommand::InsertInPlaylist {
2162 items: inserted,
2163 after: a_id,
2164 });
2165 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2166
2167 player.process_command(PlayerCommand::Undo);
2168 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2169 }
2170
2171 #[test]
2174 fn undo_move_restores_position() {
2175 let mut player = Player::new();
2176 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2177 let a_id = items[0].id;
2178 let c_id = items[2].id;
2179
2180 player.process_command(PlayerCommand::AddToPlaylist(items));
2181
2182 player.process_command(PlayerCommand::MoveInPlaylist {
2184 id: a_id,
2185 target: c_id,
2186 after: true,
2187 });
2188 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2189
2190 player.process_command(PlayerCommand::Undo);
2191 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2192 }
2193
2194 #[test]
2195 fn redo_move() {
2196 let mut player = Player::new();
2197 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2198 let a_id = items[0].id;
2199 let c_id = items[2].id;
2200
2201 player.process_command(PlayerCommand::AddToPlaylist(items));
2202 player.process_command(PlayerCommand::MoveInPlaylist {
2203 id: a_id,
2204 target: c_id,
2205 after: true,
2206 });
2207 player.process_command(PlayerCommand::Undo);
2208 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2209
2210 player.process_command(PlayerCommand::Redo);
2211 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2212 }
2213
2214 #[test]
2217 fn undo_batch_move() {
2218 let mut player = Player::new();
2219 let items = vec![
2220 make_item("A"),
2221 make_item("B"),
2222 make_item("C"),
2223 make_item("D"),
2224 ];
2225 let a_id = items[0].id;
2226 let b_id = items[1].id;
2227 let d_id = items[3].id;
2228
2229 player.process_command(PlayerCommand::AddToPlaylist(items));
2230
2231 player.process_command(PlayerCommand::MoveItemsInPlaylist {
2233 ids: vec![a_id, b_id],
2234 target: d_id,
2235 after: true,
2236 });
2237 assert_eq!(playlist_titles(&player), vec!["C", "D", "A", "B"]);
2238
2239 player.process_command(PlayerCommand::Undo);
2240 assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
2241 }
2242
2243 #[test]
2246 fn undo_clear_restores_playlist() {
2247 let mut player = Player::new();
2248 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2249
2250 player.process_command(PlayerCommand::AddToPlaylist(items));
2251 player.process_command(PlayerCommand::ClearPlaylist);
2252 assert!(playlist_ids(&player).is_empty());
2253
2254 player.process_command(PlayerCommand::Undo);
2255 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2256 }
2257
2258 #[test]
2263 fn undoing_a_replace_does_not_leave_the_engine_on_an_orphaned_track() {
2264 let mut player = Player::new();
2265 let original = seed(&mut player, 3);
2266 player.shared_state.set_cursor(Some(original[0]));
2267 pretend_playing(&mut player, original[0]);
2268
2269 let replacement = vec![make_item("something else")];
2270 let orphan = replacement[0].id;
2271 player.process_command(PlayerCommand::ReplacePlaylist {
2272 items: replacement,
2273 start: 0,
2274 });
2275 pretend_playing(&mut player, orphan);
2277
2278 player.process_command(PlayerCommand::Undo);
2279
2280 assert_eq!(playlist_ids(&player), original, "the queue comes back");
2281 assert!(
2282 player.shared_state.get_item(orphan).is_none(),
2283 "and the replacement is gone from it"
2284 );
2285 assert!(
2286 playing_id(&player).is_none_or(|id| player.shared_state.get_item(id).is_some()),
2287 "so nothing may still be playing out of it"
2288 );
2289 }
2290
2291 #[test]
2293 fn undoing_an_add_does_not_leave_the_engine_on_a_removed_track() {
2294 let mut player = Player::new();
2295 seed(&mut player, 2);
2296 let added = seed(&mut player, 1);
2297 pretend_playing(&mut player, added[0]);
2298
2299 player.process_command(PlayerCommand::Undo);
2300
2301 assert!(player.shared_state.get_item(added[0]).is_none());
2302 assert!(
2303 playing_id(&player).is_none_or(|id| player.shared_state.get_item(id).is_some()),
2304 "the engine cannot be left on the item the undo removed"
2305 );
2306 }
2307
2308 #[test]
2310 fn undoing_a_move_leaves_playback_alone() {
2311 let mut player = Player::new();
2312 let ids = seed(&mut player, 3);
2313 player.shared_state.set_cursor(Some(ids[0]));
2314 pretend_playing(&mut player, ids[0]);
2315 let starts = player.playback_starts;
2316
2317 player.process_command(PlayerCommand::MoveInPlaylist {
2318 id: ids[2],
2319 target: ids[0],
2320 after: false,
2321 });
2322 player.process_command(PlayerCommand::Undo);
2323
2324 assert_eq!(playlist_ids(&player), ids);
2325 assert_eq!(playing_id(&player), Some(ids[0]), "still on the same track");
2326 assert_eq!(player.playback_starts, starts, "and not restarted");
2327 }
2328
2329 #[test]
2330 fn redo_clear() {
2331 let mut player = Player::new();
2332 let items = vec![make_item("A"), make_item("B")];
2333
2334 player.process_command(PlayerCommand::AddToPlaylist(items));
2335 player.process_command(PlayerCommand::ClearPlaylist);
2336 player.process_command(PlayerCommand::Undo);
2337 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2338
2339 player.process_command(PlayerCommand::Redo);
2340 assert!(playlist_ids(&player).is_empty());
2341 }
2342
2343 #[test]
2346 fn multiple_undos_in_sequence() {
2347 let mut player = Player::new();
2348
2349 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("A")]));
2350 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
2351 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("C")]));
2352 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2353
2354 player.process_command(PlayerCommand::Undo);
2355 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2356
2357 player.process_command(PlayerCommand::Undo);
2358 assert_eq!(playlist_titles(&player), vec!["A"]);
2359
2360 player.process_command(PlayerCommand::Undo);
2361 assert!(playlist_ids(&player).is_empty());
2362 }
2363
2364 #[test]
2365 fn undo_redo_undo_cycle() {
2366 let mut player = Player::new();
2367 let items = vec![make_item("A"), make_item("B")];
2368
2369 player.process_command(PlayerCommand::AddToPlaylist(items));
2370 player.process_command(PlayerCommand::Undo);
2371 assert!(playlist_ids(&player).is_empty());
2372
2373 player.process_command(PlayerCommand::Redo);
2374 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2375
2376 player.process_command(PlayerCommand::Undo);
2377 assert!(playlist_ids(&player).is_empty());
2378 }
2379
2380 #[test]
2381 fn new_action_clears_redo_stack() {
2382 let mut player = Player::new();
2383 let items = vec![make_item("A")];
2384
2385 player.process_command(PlayerCommand::AddToPlaylist(items));
2386 player.process_command(PlayerCommand::Undo);
2387 assert!(player.undo_stack().can_redo());
2388
2389 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
2391 assert!(!player.undo_stack().can_redo());
2392 }
2393
2394 #[test]
2395 fn undo_on_empty_stack_is_noop() {
2396 let mut player = Player::new();
2397 player.process_command(PlayerCommand::Undo);
2398 assert!(playlist_ids(&player).is_empty());
2399 }
2400
2401 #[test]
2402 fn redo_on_empty_stack_is_noop() {
2403 let mut player = Player::new();
2404 player.process_command(PlayerCommand::Redo);
2405 assert!(playlist_ids(&player).is_empty());
2406 }
2407
2408 #[test]
2411 fn playback_commands_not_undoable() {
2412 let mut player = Player::new();
2413 player.process_command(PlayerCommand::Pause);
2414 player.process_command(PlayerCommand::Resume);
2415 player.process_command(PlayerCommand::NextTrack);
2416 player.process_command(PlayerCommand::PrevTrack);
2417 assert!(!player.undo_stack().can_undo());
2418 }
2419
2420 #[test]
2421 fn update_paths_not_undoable() {
2422 let mut player = Player::new();
2423 let items = vec![make_item("A")];
2424 let id = items[0].id;
2425 player.process_command(PlayerCommand::AddToPlaylist(items));
2426
2427 let undo_count = player.undo_stack().undo_len();
2428 player.process_command(PlayerCommand::UpdatePaths(vec![(
2429 id,
2430 PathBuf::from("/new/path.flac"),
2431 )]));
2432 assert_eq!(player.undo_stack().undo_len(), undo_count);
2433 }
2434
2435 #[test]
2438 fn add_remove_undo_undo_produces_original() {
2439 let mut player = Player::new();
2440 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2441 let b_id = items[1].id;
2442 let original_titles = vec!["A", "B", "C"];
2443
2444 player.process_command(PlayerCommand::AddToPlaylist(items));
2445 player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
2446 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2447
2448 player.process_command(PlayerCommand::Undo);
2450 assert_eq!(playlist_titles(&player), original_titles);
2451
2452 player.process_command(PlayerCommand::Undo);
2454 assert!(playlist_ids(&player).is_empty());
2455 }
2456
2457 #[test]
2458 fn interleaved_adds_and_moves_undo() {
2459 let mut player = Player::new();
2460 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2461 let a_id = items[0].id;
2462 let c_id = items[2].id;
2463
2464 player.process_command(PlayerCommand::AddToPlaylist(items));
2465
2466 player.process_command(PlayerCommand::MoveInPlaylist {
2468 id: a_id,
2469 target: c_id,
2470 after: true,
2471 });
2472 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2473
2474 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("D")]));
2476 assert_eq!(playlist_titles(&player), vec!["B", "C", "A", "D"]);
2477
2478 player.process_command(PlayerCommand::Undo);
2480 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2481
2482 player.process_command(PlayerCommand::Undo);
2484 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2485 }
2486
2487 #[test]
2492 fn stop_engine_drops_engine_synchronously() {
2493 use std::sync::atomic::{AtomicBool, Ordering};
2494
2495 struct MockEngine {
2496 dropped: Arc<AtomicBool>,
2497 }
2498 impl AudioEngineHandle for MockEngine {
2499 fn start(&self) -> Result<(), BackendError> {
2500 Ok(())
2501 }
2502 fn stop(&self) -> Result<(), BackendError> {
2503 Ok(())
2504 }
2505 fn is_running(&self) -> bool {
2506 false
2507 }
2508 fn fade_out(&self) {}
2509 fn fade_in(&self) -> Result<(), BackendError> {
2510 Ok(())
2511 }
2512 fn is_silent(&self) -> bool {
2513 false
2514 }
2515 }
2516 impl Drop for MockEngine {
2517 fn drop(&mut self) {
2518 self.dropped.store(true, Ordering::SeqCst);
2519 }
2520 }
2521
2522 let dropped = Arc::new(AtomicBool::new(false));
2523
2524 let stop_flag = Arc::new(AtomicBool::new(false));
2526 let decode_handle = buffer::DecodeHandle::new_for_test(stop_flag);
2527
2528 let mut player = Player::new();
2529 player.active_playback = Some(ActivePlayback {
2530 engine: Box::new(MockEngine {
2531 dropped: dropped.clone(),
2532 }),
2533 decode_handle,
2534 stream: None,
2535 _rate_watch: None,
2536 });
2537
2538 player.stop_engine();
2539
2540 assert!(
2544 dropped.load(Ordering::SeqCst),
2545 "AudioEngine must be dropped synchronously in stop_engine (GitHub #89)"
2546 );
2547 }
2548
2549 #[test]
2550 fn abandoning_a_stream_wakes_a_reader_parked_at_the_write_head() {
2551 let live = LiveStream {
2552 feed: crate::remote::downloads::ByteFeed::new(),
2553 abandoned: Default::default(),
2554 };
2555 let feed = live.feed.clone();
2556 let started = std::time::Instant::now();
2557 let reader = thread::spawn(move || {
2558 feed.wait_past(
2559 0,
2560 std::time::Instant::now() + std::time::Duration::from_secs(30),
2561 )
2562 });
2563 thread::sleep(std::time::Duration::from_millis(50));
2564 live.abandon();
2565 reader.join().unwrap();
2566
2567 assert!(live.abandoned.load(std::sync::atomic::Ordering::Acquire));
2568 assert!(started.elapsed() < std::time::Duration::from_secs(5));
2569 }
2570
2571 struct StuckBackend {
2576 rate: f64,
2577 asked: Arc<std::sync::Mutex<Option<(f64, u32)>>>,
2578 }
2579
2580 struct NullEngine;
2581 impl AudioEngineHandle for NullEngine {
2582 fn start(&self) -> Result<(), BackendError> {
2583 Ok(())
2584 }
2585 fn stop(&self) -> Result<(), BackendError> {
2586 Ok(())
2587 }
2588 fn is_running(&self) -> bool {
2589 false
2590 }
2591 fn fade_out(&self) {}
2592 fn fade_in(&self) -> Result<(), BackendError> {
2593 Ok(())
2594 }
2595 fn is_silent(&self) -> bool {
2596 false
2597 }
2598 }
2599
2600 impl AudioBackend for StuckBackend {
2601 fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
2602 Ok(vec![self.default_device()?])
2603 }
2604 fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
2605 Ok(backend::DeviceInfo {
2606 name: "Stuck DAC".into(),
2607 sample_rates: vec![self.rate],
2608 platform_id: 0,
2609 })
2610 }
2611 fn supported_sample_rates(
2612 &self,
2613 _device: &backend::DeviceInfo,
2614 ) -> Result<Vec<f64>, BackendError> {
2615 Ok(vec![self.rate])
2616 }
2617 fn get_device_sample_rate(
2618 &self,
2619 _device: &backend::DeviceInfo,
2620 ) -> Result<f64, BackendError> {
2621 Ok(self.rate)
2622 }
2623 fn set_device_sample_rate(
2624 &self,
2625 _device: &backend::DeviceInfo,
2626 rate: f64,
2627 ) -> Result<f64, BackendError> {
2628 Err(BackendError::UnsupportedSampleRate(rate))
2629 }
2630 fn create_engine(
2631 &self,
2632 _device: &backend::DeviceInfo,
2633 sample_rate: f64,
2634 channels: u32,
2635 _consumer: rtrb::Consumer<f32>,
2636 _samples_played: Arc<AtomicU64>,
2637 ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
2638 *self.asked.lock().unwrap() = Some((sample_rate, channels));
2639 Ok(Box::new(NullEngine))
2640 }
2641 }
2642
2643 fn engine_format_for(source_rate: u32, channels: u16, device_rate: f64) -> (f64, u32) {
2644 let asked = Arc::new(std::sync::Mutex::new(None));
2645 let mut player = Player::new();
2646 player.backend = Box::new(StuckBackend {
2647 rate: device_rate,
2648 asked: asked.clone(),
2649 });
2650
2651 let info = buffer::StreamInfo {
2652 codec: "MP3".into(),
2653 sample_rate: source_rate,
2654 channels,
2655 bit_depth: Some(16),
2656 bitrate_kbps: None,
2657 duration_ms: 1000,
2658 };
2659 let (_producer, consumer) = rtrb::RingBuffer::new(16);
2660 player
2661 .create_engine_for(&info, consumer)
2662 .expect("engine creation should succeed");
2663 let asked = *asked.lock().unwrap();
2664 asked.expect("engine was never created")
2665 }
2666
2667 #[test]
2668 fn engine_uses_source_rate_when_device_refuses_switch() {
2669 assert_eq!(engine_format_for(22050, 2, 48000.0), (22050.0, 2));
2672 assert_eq!(engine_format_for(32000, 2, 44100.0), (32000.0, 2));
2673 }
2674
2675 #[test]
2676 fn engine_uses_source_channel_count() {
2677 assert_eq!(engine_format_for(44100, 1, 44100.0), (44100.0, 1));
2678 }
2679
2680 fn settled_rate_for(source_rate: u32, device_rate: f64) -> Option<u32> {
2682 let mut player = Player::new();
2683 player.backend = Box::new(StuckBackend {
2684 rate: device_rate,
2685 asked: Arc::new(std::sync::Mutex::new(None)),
2686 });
2687 let state = player.shared_state.clone();
2688
2689 let info = buffer::StreamInfo {
2690 codec: "MP3".into(),
2691 sample_rate: source_rate,
2692 channels: 2,
2693 bit_depth: Some(16),
2694 bitrate_kbps: None,
2695 duration_ms: 1000,
2696 };
2697 let (_producer, consumer) = rtrb::RingBuffer::new(16);
2698 player
2699 .create_engine_for(&info, consumer)
2700 .expect("engine creation should succeed");
2701 state.output_sample_rate()
2702 }
2703
2704 #[test]
2705 fn settled_device_rate_reaches_the_shared_state() {
2706 assert_eq!(settled_rate_for(22050, 48000.0), Some(48000));
2710 assert_eq!(settled_rate_for(44100, 44100.0), Some(44100));
2712 }
2713
2714 struct SlowBackend {
2716 observed: Arc<std::sync::Mutex<Vec<Option<u32>>>>,
2717 state: Arc<SharedPlayerState>,
2718 }
2719
2720 impl AudioBackend for SlowBackend {
2721 fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
2722 Ok(vec![self.default_device()?])
2723 }
2724 fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
2725 Ok(backend::DeviceInfo {
2726 name: "Slow DAC".into(),
2727 sample_rates: vec![44100.0, 48000.0],
2728 platform_id: 0,
2729 })
2730 }
2731 fn supported_sample_rates(
2732 &self,
2733 _device: &backend::DeviceInfo,
2734 ) -> Result<Vec<f64>, BackendError> {
2735 Ok(vec![44100.0, 48000.0])
2736 }
2737 fn get_device_sample_rate(
2738 &self,
2739 _device: &backend::DeviceInfo,
2740 ) -> Result<f64, BackendError> {
2741 Ok(48000.0)
2742 }
2743 fn set_device_sample_rate(
2744 &self,
2745 _device: &backend::DeviceInfo,
2746 rate: f64,
2747 ) -> Result<f64, BackendError> {
2748 self.observed
2750 .lock()
2751 .unwrap()
2752 .push(self.state.output_sample_rate());
2753 Ok(rate)
2754 }
2755 fn create_engine(
2756 &self,
2757 _device: &backend::DeviceInfo,
2758 _sample_rate: f64,
2759 _channels: u32,
2760 _consumer: rtrb::Consumer<f32>,
2761 _samples_played: Arc<AtomicU64>,
2762 ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
2763 Ok(Box::new(NullEngine))
2764 }
2765 }
2766
2767 #[test]
2768 fn the_previous_rate_is_not_published_while_the_device_reclocks() {
2769 let mut player = Player::new();
2775 let state = player.shared_state.clone();
2776 state.set_output_sample_rate(48000);
2777
2778 let observed = Arc::new(std::sync::Mutex::new(Vec::new()));
2779 player.backend = Box::new(SlowBackend {
2780 observed: observed.clone(),
2781 state: state.clone(),
2782 });
2783
2784 let info = buffer::StreamInfo {
2785 codec: "FLAC".into(),
2786 sample_rate: 44100,
2787 channels: 2,
2788 bit_depth: Some(16),
2789 bitrate_kbps: None,
2790 duration_ms: 1000,
2791 };
2792 let (_producer, consumer) = rtrb::RingBuffer::new(16);
2793 player
2794 .create_engine_for(&info, consumer)
2795 .expect("engine creation should succeed");
2796
2797 assert_eq!(
2798 *observed.lock().unwrap(),
2799 vec![None],
2800 "mid-switch the output rate must read as unknown, not as the last track's"
2801 );
2802 assert_eq!(state.output_sample_rate(), Some(44100));
2803 }
2804
2805 struct WatchedBackend {
2807 inner: StuckBackend,
2808 #[allow(clippy::type_complexity)]
2809 captured: Arc<std::sync::Mutex<Option<Box<dyn Fn(f64) + Send + Sync>>>>,
2810 }
2811
2812 struct NullWatch;
2813 impl backend::SampleRateWatch for NullWatch {}
2814
2815 impl AudioBackend for WatchedBackend {
2816 fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
2817 self.inner.list_devices()
2818 }
2819 fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
2820 self.inner.default_device()
2821 }
2822 fn supported_sample_rates(
2823 &self,
2824 device: &backend::DeviceInfo,
2825 ) -> Result<Vec<f64>, BackendError> {
2826 self.inner.supported_sample_rates(device)
2827 }
2828 fn get_device_sample_rate(
2829 &self,
2830 device: &backend::DeviceInfo,
2831 ) -> Result<f64, BackendError> {
2832 self.inner.get_device_sample_rate(device)
2833 }
2834 fn set_device_sample_rate(
2835 &self,
2836 device: &backend::DeviceInfo,
2837 rate: f64,
2838 ) -> Result<f64, BackendError> {
2839 self.inner.set_device_sample_rate(device, rate)
2840 }
2841 fn watch_device_sample_rate(
2842 &self,
2843 _device: &backend::DeviceInfo,
2844 on_change: Box<dyn Fn(f64) + Send + Sync>,
2845 ) -> Option<Box<dyn backend::SampleRateWatch>> {
2846 *self.captured.lock().unwrap() = Some(on_change);
2847 Some(Box::new(NullWatch))
2848 }
2849 fn create_engine(
2850 &self,
2851 device: &backend::DeviceInfo,
2852 sample_rate: f64,
2853 channels: u32,
2854 consumer: rtrb::Consumer<f32>,
2855 samples_played: Arc<AtomicU64>,
2856 ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
2857 self.inner
2858 .create_engine(device, sample_rate, channels, consumer, samples_played)
2859 }
2860 }
2861
2862 #[test]
2863 fn external_rate_change_reaches_the_shared_state() {
2864 let captured = Arc::new(std::sync::Mutex::new(None));
2868 let mut player = Player::new();
2869 player.backend = Box::new(WatchedBackend {
2870 inner: StuckBackend {
2871 rate: 44100.0,
2872 asked: Arc::new(std::sync::Mutex::new(None)),
2873 },
2874 captured: captured.clone(),
2875 });
2876 let state = player.shared_state.clone();
2877
2878 let info = buffer::StreamInfo {
2879 codec: "FLAC".into(),
2880 sample_rate: 44100,
2881 channels: 2,
2882 bit_depth: Some(16),
2883 bitrate_kbps: None,
2884 duration_ms: 1000,
2885 };
2886 let (_producer, consumer) = rtrb::RingBuffer::new(16);
2887 player
2888 .create_engine_for(&info, consumer)
2889 .expect("engine creation should succeed");
2890 assert_eq!(state.output_sample_rate(), Some(44100));
2891
2892 let on_change = captured.lock().unwrap().take().expect("watch registered");
2893 on_change(48000.0);
2894 assert_eq!(state.output_sample_rate(), Some(48000));
2895 }
2896}