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.clone(),
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();
463
464 let advance_state = self.shared_state.clone();
468 let decode_cursor = parking_lot::Mutex::new(Some(id));
469 let next_track = move || {
470 let current = decode_cursor.lock().take()?;
471 let next = advance_state.peek_next_ready_after(current);
472 if let Some((next_id, _)) = &next {
473 let mut guard = decode_cursor.lock();
474 *guard = Some(*next_id);
475 }
476 next
477 };
478
479 let cfg = crate::config::Config::load_or_default();
481 let rg_mode = cfg.playback.replaygain;
482 let pre_amp_db = cfg.playback.pre_amp_db;
483
484 let finish_tx = self.commands.tx.clone();
485 let (_stream_info, decode_handle) = buffer::start_decode_file(
486 id,
487 path,
488 producer,
489 seek_ms,
490 next_track,
491 self.timeline.clone(),
492 Some(self.viz_buffer.clone()),
493 rg_mode,
494 pre_amp_db,
495 move || {
496 finish_tx.send(PlayerCommand::DecodeFinished).ok();
497 },
498 )?;
499
500 let (engine, rate_watch) = self.create_engine_for(&info, consumer)?;
501 engine.start()?;
502
503 self.shared_state.set_playback_state(PlaybackState::Playing);
504
505 self.active_playback = Some(ActivePlayback {
506 engine,
507 decode_handle,
508 stream: None,
509 _rate_watch: rate_watch,
510 });
511
512 Ok(())
513 }
514
515 fn probe_stream_for_playback(
530 &self,
531 id: QueueItemId,
532 path: &Path,
533 bytes_written: Arc<crate::remote::downloads::ByteFeed>,
534 total: u64,
535 ) {
536 let path = path.to_path_buf();
537 let tx = self.commands.tx.clone();
538 let hint = hint_for(&path);
539
540 let status = {
545 let downloading = self.stream_status_fn(id);
546 let state = self.shared_state.clone();
547 Arc::new(move || {
548 if state.is_cursor(id) {
549 downloading()
550 } else {
551 streaming::StreamStatus::Failed
552 }
553 }) as Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync>
554 };
555
556 let spawned = thread::Builder::new()
557 .name("koan-stream-probe".into())
558 .spawn(move || {
559 let attempt = |mode, wait: bool| {
566 let open = if wait {
567 streaming::PartialFileSource::open(
568 &path,
569 bytes_written.clone(),
570 total,
571 status.clone(),
572 mode,
573 )
574 } else {
575 streaming::PartialFileSource::open_for_probe(
576 &path,
577 bytes_written.clone(),
578 total,
579 status.clone(),
580 mode,
581 )
582 };
583 open.map_err(buffer::DecodeError::Io).and_then(|source| {
584 let mss = symphonia::core::io::MediaSourceStream::new(
585 Box::new(source),
586 Default::default(),
587 );
588 buffer::probe_source(mss, &hint)
589 })
590 };
591
592 let info = match attempt(streaming::ProbeMode::Full, false) {
596 Ok(info) => Some((info, streaming::ProbeMode::Full)),
597 Err(e) => {
598 log::info!(
604 "stream probe: {} needs more than has arrived ({}), opening without a length",
605 path.display(),
606 e
607 );
608 let lengthless = lengthless_mode_for(&path);
609 attempt(lengthless, true)
610 .ok()
611 .map(|info| (info, lengthless))
612 }
613 };
614
615 match info {
616 Some((info, mode)) => {
617 tx.send(PlayerCommand::StreamProbed {
618 id,
619 info: Box::new(info),
620 mode,
621 })
622 .ok();
623 }
624 None => log::info!(
627 "stream probe: {} cannot start early, waiting for the download",
628 path.display()
629 ),
630 }
631 });
632
633 if let Err(e) = spawned {
634 log::warn!("stream probe: could not spawn for {:?}: {}", id, e);
635 }
636 }
637
638 fn stream_probed(
641 &mut self,
642 id: QueueItemId,
643 info: buffer::StreamInfo,
644 mode: streaming::ProbeMode,
645 ) {
646 if !self.shared_state.is_cursor(id) {
647 return; }
649 if self.shared_state.playback_state() != PlaybackState::Stopped {
650 return; }
652
653 match self.shared_state.item_playback_source(id) {
654 Some(PlaybackSource::Ready(path)) => {
656 if let Err(e) = self.start_playback(id, &path, 0) {
657 log::error!("stream probe: playback failed: {}", e);
658 }
659 }
660 Some(PlaybackSource::Streaming {
661 path,
662 bytes_written,
663 total,
664 }) => {
665 let source = StreamSource {
666 path,
667 bytes_written,
668 total,
669 mode,
670 };
671 if let Err(e) = self.start_streaming_playback(id, source, 0, info) {
672 log::error!("stream probe: streaming playback failed: {}", e);
673 }
674 }
675 None => {}
676 }
677 }
678
679 fn stream_status_fn(
683 &self,
684 id: QueueItemId,
685 ) -> Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync> {
686 let state = self.shared_state.clone();
687 Arc::new(move || match state.item_load_state(id) {
688 Some(LoadState::Ready) => streaming::StreamStatus::Complete,
689 Some(LoadState::Failed(_)) => streaming::StreamStatus::Failed,
690 _ => streaming::StreamStatus::Downloading,
691 })
692 }
693
694 fn start_streaming_playback(
704 &mut self,
705 id: QueueItemId,
706 source: StreamSource,
707 seek_ms: u64,
708 info: buffer::StreamInfo,
709 ) -> Result<(), PlayerError> {
710 let result = self.open_streaming_playback(id, source, seek_ms, info);
711 if result.is_err() {
712 self.stop_playback_and_clear_state();
713 }
714 result
715 }
716
717 fn open_streaming_playback(
718 &mut self,
719 id: QueueItemId,
720 source: StreamSource,
721 seek_ms: u64,
722 info: buffer::StreamInfo,
723 ) -> Result<(), PlayerError> {
724 self.stop_engine();
725 self.stream_mode = source.mode;
727 let path = source.path.as_path();
728
729 let live = LiveStream {
730 feed: source.bytes_written.clone(),
731 abandoned: Default::default(),
732 };
733 let status = {
734 let downloading = self.stream_status_fn(id);
735 let abandoned = live.abandoned.clone();
736 Arc::new(move || {
737 if abandoned.load(std::sync::atomic::Ordering::Acquire) {
738 streaming::StreamStatus::Failed
739 } else {
740 downloading()
741 }
742 }) as Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync>
743 };
744 let open_source = {
745 let StreamSource {
746 path,
747 bytes_written,
748 total,
749 mode,
750 } = source.clone();
751 let status = status.clone();
752 move || {
753 streaming::PartialFileSource::open(
754 &path,
755 bytes_written.clone(),
756 total,
757 status.clone(),
758 mode,
759 )
760 }
761 };
762
763 self.shared_state.set_track_info(Some(TrackInfo {
764 id,
765 path: path.to_path_buf(),
766 codec: info.codec.clone(),
767 sample_rate: info.sample_rate,
768 bit_depth: info.bit_depth,
769 bitrate_kbps: info.bitrate_kbps,
770 channels: info.channels,
771 duration_ms: info.duration_ms,
772 }));
773 self.shared_state.set_position_ms(seek_ms);
774 self.on_track_changed(id, seek_ms);
775 log::info!(
776 "streaming: {} ({:?}) — {} {}Hz/{}ch, {}ms{}",
777 path.display(),
778 id,
779 info.codec,
780 info.sample_rate,
781 info.channels,
782 info.duration_ms,
783 if seek_ms > 0 {
784 format!(" @{}ms", seek_ms)
785 } else {
786 String::new()
787 },
788 );
789
790 let (producer, consumer) = rtrb::RingBuffer::new(RING_BUFFER_SIZE);
791
792 self.timeline.reset();
793
794 let advance_state = self.shared_state.clone();
796 let decode_cursor = parking_lot::Mutex::new(Some(id));
797 let next_track = move || {
798 let current = decode_cursor.lock().take()?;
799 let next = advance_state.peek_next_ready_after(current);
800 if let Some((next_id, _)) = &next {
801 let mut guard = decode_cursor.lock();
802 *guard = Some(*next_id);
803 }
804 next
805 };
806
807 let first = buffer::SourceEntry {
808 id,
809 path: path.to_path_buf(),
810 hint: hint_for(path),
811 make_mss: Box::new(move || {
812 Ok(symphonia::core::io::MediaSourceStream::new(
813 Box::new(open_source()?),
814 Default::default(),
815 ))
816 }),
817 };
818
819 let cfg = crate::config::Config::load_or_default();
821 let rg_mode = cfg.playback.replaygain;
822 let pre_amp_db = cfg.playback.pre_amp_db;
823
824 let finish_tx = self.commands.tx.clone();
825 let (_stream_info, decode_handle) = buffer::start_decode(
826 first,
827 producer,
828 seek_ms,
829 move || {
830 let (next_id, next_path) = next_track()?;
831 Some(buffer::SourceEntry::from_file(next_id, next_path))
832 },
833 self.timeline.clone(),
834 Some(self.viz_buffer.clone()),
835 rg_mode,
836 pre_amp_db,
837 move || {
838 finish_tx.send(PlayerCommand::DecodeFinished).ok();
839 },
840 )?;
841
842 let (engine, rate_watch) = self.create_engine_for(&info, consumer)?;
843 engine.start()?;
844
845 self.shared_state.set_playback_state(PlaybackState::Playing);
846
847 self.active_playback = Some(ActivePlayback {
848 engine,
849 decode_handle,
850 stream: Some(live),
851 _rate_watch: rate_watch,
852 });
853
854 Ok(())
855 }
856
857 pub fn seek(&mut self, position_ms: u64) {
864 let Some(info) = self.shared_state.track_info() else {
865 return;
866 };
867 let seekable = self.shared_state.seekable_ms();
869 if seekable == 0 {
870 log::debug!("seek declined: {:?} is not seekable yet", info.id);
874 return;
875 }
876 let ceiling = seekable.min(
877 self.shared_state
878 .duration_ms()
879 .saturating_sub(SEEK_END_GUARD_MS),
880 );
881 let clamped = position_ms.min(ceiling);
882
883 if let Err(e) = self.restart_current(&info, clamped) {
884 log::error!("seek failed: {}", e);
885 }
886 }
887
888 fn restart_current(&mut self, info: &TrackInfo, position_ms: u64) -> Result<(), PlayerError> {
896 let was_paused = self.shared_state.playback_state() == PlaybackState::Paused;
897
898 match self.shared_state.item_playback_source(info.id) {
899 Some(PlaybackSource::Streaming {
900 path,
901 bytes_written,
902 total,
903 }) => {
904 let known = buffer::StreamInfo {
908 codec: info.codec.clone(),
909 sample_rate: info.sample_rate,
910 channels: info.channels,
911 bit_depth: info.bit_depth,
912 bitrate_kbps: info.bitrate_kbps,
913 duration_ms: info.duration_ms,
914 };
915 let source = StreamSource {
916 path,
917 bytes_written,
918 total,
919 mode: self.stream_mode,
920 };
921 self.start_streaming_playback(info.id, source, position_ms, known)?;
922 }
923 Some(PlaybackSource::Ready(path)) => {
924 self.start_playback(info.id, &path, position_ms)?;
925 }
926 None => return Ok(()),
927 }
928
929 if was_paused {
930 self.pause_now();
931 }
932 self.report(if was_paused {
933 PlaybackReportState::Paused
934 } else {
935 PlaybackReportState::Playing
936 });
937 Ok(())
938 }
939
940 pub fn next_track(&mut self) {
942 match self.shared_state.advance_cursor_loadable() {
943 Some(id) => self.play(id),
944 None => {
945 log::info!("no more tracks in playlist");
946 self.stop_playback_and_clear_state();
947 }
948 }
949 }
950
951 pub fn prev_track(&mut self) {
953 match self.shared_state.retreat_cursor() {
954 Some((id, _)) => self.play(id),
955 None => {
956 if let Some(info) = self.shared_state.track_info()
958 && let Err(e) = self.restart_current(&info, 0)
959 {
960 log::error!("restart failed: {}", e);
961 }
962 }
963 }
964 }
965
966 pub fn pause(&mut self) {
971 let Some(ref playback) = self.active_playback else {
972 return;
973 };
974 if crate::config::Config::load_or_default()
975 .playback
976 .fade_on_pause
977 {
978 playback.engine.fade_out();
979 self.shared_state.set_playback_state(PlaybackState::Paused);
980 } else {
981 self.pause_now();
982 }
983 self.report(PlaybackReportState::Paused);
984 }
985
986 fn pause_now(&mut self) {
989 if let Some(ref playback) = self.active_playback {
990 if let Err(e) = playback.engine.stop() {
991 log::error!("pause failed: {}", e);
992 return;
993 }
994 self.shared_state.set_playback_state(PlaybackState::Paused);
995 }
996 }
997
998 pub fn resume(&mut self) {
1000 if let Some(ref playback) = self.active_playback {
1001 let engine = &playback.engine;
1002 let resumed = if engine.is_running() || engine.is_silent() {
1003 engine.fade_in()
1004 } else {
1005 engine.start()
1006 };
1007 if let Err(e) = resumed {
1008 log::error!("resume failed: {}", e);
1009 return;
1010 }
1011 self.shared_state.set_playback_state(PlaybackState::Playing);
1012 self.wake_analyzer();
1013 self.report(PlaybackReportState::Playing);
1014 }
1015 }
1016
1017 fn wake_analyzer(&self) {
1024 self.viz_snapshot.wake();
1025 }
1026
1027 pub fn stop(&mut self) {
1029 self.shared_state.clear_playlist();
1030 self.stop_playback_and_clear_state();
1031 }
1032
1033 fn stop_engine(&mut self) {
1039 let Some(playback) = self.active_playback.take() else {
1040 return;
1041 };
1042 let ActivePlayback {
1043 engine,
1044 mut decode_handle,
1045 stream,
1046 _rate_watch,
1047 } = playback;
1048
1049 let _ = engine.stop();
1050 decode_handle.signal_stop();
1053 if let Some(stream) = stream {
1054 stream.abandon();
1055 }
1056 decode_handle.stop();
1057 drop(engine);
1058 }
1059
1060 fn stop_playback_and_clear_state(&mut self) {
1062 self.report(PlaybackReportState::Stopped);
1063 self.finish_play();
1064 self.stop_engine();
1065 self.timeline.reset();
1066 self.shared_state.set_playback_state(PlaybackState::Stopped);
1067 self.shared_state.set_position_ms(0);
1068 self.shared_state.set_track_info(None);
1069 }
1070
1071 pub fn remove_from_playlist(&mut self, id: QueueItemId) {
1078 let was_cursor = self.shared_state.is_cursor(id);
1079 let resume_after = was_cursor
1080 .then(|| self.shared_state.item_before(id))
1081 .flatten();
1082 self.shared_state.remove_item(id);
1083 if was_cursor {
1084 self.shared_state.set_cursor(resume_after);
1085 self.next_track();
1086 }
1087 }
1088
1089 pub fn track_ready(&mut self, id: QueueItemId) {
1092 self.shared_state.update_item_state(id, ItemState::Ready);
1094
1095 if !self.shared_state.is_cursor(id) {
1096 return;
1097 }
1098
1099 let is_playing = self.shared_state.playback_state() == PlaybackState::Playing;
1100 let current_track_id = self.shared_state.track_info().map(|t| t.id);
1101
1102 if is_playing && current_track_id == Some(id) {
1103 log::info!(
1106 "track_ready: download complete while streaming {:?}, refreshing metadata",
1107 id
1108 );
1109 self.refresh_track_metadata(id);
1110 return;
1111 }
1112
1113 if !is_playing && let Some(path) = self.shared_state.item_path_if_ready(id) {
1115 log::info!("track_ready: starting playback for {:?}", id);
1116 if let Err(e) = self.start_playback(id, &path, 0) {
1117 log::error!("track_ready playback failed: {}", e);
1118 }
1119 }
1120 }
1121
1122 pub fn track_stream_ready(&mut self, id: QueueItemId) {
1125 if !self.shared_state.is_cursor(id) {
1126 return;
1127 }
1128
1129 let is_playing = self.shared_state.playback_state() == PlaybackState::Playing;
1130 if is_playing {
1131 return; }
1133
1134 match self.shared_state.item_playback_source(id) {
1135 Some(PlaybackSource::Streaming {
1136 path,
1137 bytes_written,
1138 total,
1139 }) => {
1140 log::info!("track_stream_ready: probing partial file for {:?}", id);
1141 self.probe_stream_for_playback(id, &path, bytes_written, total);
1142 }
1143 Some(PlaybackSource::Ready(path)) => {
1144 log::info!(
1146 "track_stream_ready: track already ready, starting normal playback for {:?}",
1147 id
1148 );
1149 if let Err(e) = self.start_playback(id, &path, 0) {
1150 log::error!("track_stream_ready playback failed: {}", e);
1151 }
1152 }
1153 None => {} }
1155 }
1156
1157 fn refresh_track_metadata(&mut self, id: QueueItemId) {
1161 use crate::index::metadata;
1162
1163 let path = match self.shared_state.item_path_if_ready(id) {
1164 Some(p) => p,
1165 None => return,
1166 };
1167
1168 match metadata::read_metadata(&path) {
1169 Ok(meta) => {
1170 self.shared_state.update_item_metadata(
1171 id,
1172 meta.title,
1173 meta.artist,
1174 meta.album_artist.unwrap_or_default(),
1175 meta.album,
1176 meta.duration_ms.map(|d| d as u64),
1177 );
1178
1179 if let Some(current) = self.shared_state.track_info()
1188 && current.id == id
1189 {
1190 let probed = buffer::probe_file(&path).ok();
1191 let duration_ms = probed
1192 .as_ref()
1193 .map(|s| s.duration_ms)
1194 .filter(|d| *d > current.duration_ms)
1195 .unwrap_or(current.duration_ms);
1196 if duration_ms != current.duration_ms {
1197 log::info!(
1198 "track_ready: duration corrected {}ms → {}ms",
1199 current.duration_ms,
1200 duration_ms
1201 );
1202 }
1203 self.shared_state.set_track_info(Some(TrackInfo {
1204 duration_ms,
1205 path: path.clone(),
1206 ..current
1207 }));
1208 }
1209
1210 self.shared_state.signal_metadata_refresh();
1212 log::info!("track_ready: metadata refreshed for {:?}", id);
1213 }
1214 Err(e) => {
1215 log::warn!("track_ready: metadata refresh failed for {:?}: {}", id, e);
1216 }
1217 }
1218 }
1219
1220 fn on_track_changed(&mut self, id: QueueItemId, position_ms: u64) {
1230 if self.in_flight.as_ref().is_some_and(|f| f.item == id) {
1231 return;
1232 }
1233 self.finish_play();
1234 let track_id = self.shared_state.item_db_id(id);
1235 self.in_flight = Some(InFlight::new(id, track_id));
1236 if let (Some(track_id), Some(recorder)) = (track_id, self.history.as_ref()) {
1237 recorder.record(PlayEvent::Started {
1238 track_id,
1239 position_ms,
1240 });
1241 }
1242 }
1243
1244 fn report(&self, state: PlaybackReportState) {
1248 let Some(track_id) = self.in_flight.as_ref().and_then(InFlight::track_id) else {
1249 return;
1250 };
1251 if let Some(recorder) = self.history.as_ref() {
1252 recorder.record(PlayEvent::Playback(PlaybackReport {
1253 track_id,
1254 state,
1255 position_ms: self.shared_state.position_ms(),
1256 }));
1257 }
1258 }
1259
1260 fn finish_play(&mut self) -> Option<PlayEvent> {
1263 let flight = self.in_flight.take()?;
1264 let event = PlayEvent::Finished {
1265 track_id: flight.track_id()?,
1266 listened_ms: flight.listened_ms(),
1267 };
1268 if let Some(recorder) = self.history.as_ref() {
1269 recorder.record(event);
1270 }
1271 Some(event)
1272 }
1273
1274 pub fn update_playback_state(&mut self) {
1275 let Some(playback) = self.active_playback.as_ref() else {
1276 return;
1277 };
1278
1279 if self.shared_state.playback_state() == PlaybackState::Paused
1280 && playback.engine.is_running()
1281 && playback.engine.is_silent()
1282 && let Err(e) = playback.engine.stop()
1283 {
1284 log::error!("stopping after fade failed: {}", e);
1285 }
1286
1287 if let Some((id, path, info, position_ms)) = self.timeline.current_playback() {
1288 self.shared_state.set_position_ms(position_ms);
1289
1290 self.on_track_changed(id, position_ms);
1293 if let Some(f) = self.in_flight.as_mut() {
1294 f.advance(position_ms);
1295 }
1296
1297 let current_id = self.shared_state.track_info().map(|t| t.id);
1300 if current_id != Some(id) {
1301 log::info!("timeline: now playing {:?}", id);
1302 self.shared_state.set_track_info(Some(TrackInfo {
1303 id,
1304 path,
1305 codec: info.codec,
1306 sample_rate: info.sample_rate,
1307 bit_depth: info.bit_depth,
1308 bitrate_kbps: info.bitrate_kbps,
1309 channels: info.channels,
1310 duration_ms: info.duration_ms,
1311 }));
1312 self.shared_state.set_cursor(Some(id));
1313 }
1314 }
1315 }
1316
1317 pub fn track_failed(&mut self, id: QueueItemId) {
1324 if !self.shared_state.is_cursor(id) {
1325 return;
1326 }
1327 if self.shared_state.playback_state() != PlaybackState::Stopped {
1331 return;
1332 }
1333 log::info!("track {:?} cannot load, moving on", id);
1334 self.next_track();
1335 }
1336
1337 fn on_decode_finished(&mut self) {
1344 log::info!("decode finished, checking for next track");
1345 match self.shared_state.advance_cursor_loadable() {
1346 Some(id) => self.play(id),
1347 None => {
1348 log::info!("no more tracks — stopping");
1349 self.stop_playback_and_clear_state();
1350 }
1351 }
1352 }
1353
1354 fn snapshot_for_undo(
1358 &self,
1359 ids: &[QueueItemId],
1360 ) -> Vec<(Box<state::PlaylistItem>, Option<QueueItemId>)> {
1361 self.shared_state
1362 .items_before(ids)
1363 .into_iter()
1364 .filter_map(|(id, after)| Some((Box::new(self.shared_state.get_item(id)?), after)))
1365 .collect()
1366 }
1367
1368 fn push_undo(&mut self, entry: UndoEntry) {
1370 if let Some(ref mut batch) = self.batch_buffer {
1371 batch.push(entry);
1372 } else {
1373 self.undo_stack.push(entry);
1374 }
1375 }
1376
1377 pub fn process_command(&mut self, cmd: PlayerCommand) {
1379 match cmd {
1380 PlayerCommand::Play(id) => self.play(id),
1381 PlayerCommand::Pause => self.pause(),
1382 PlayerCommand::Resume => self.resume(),
1383 PlayerCommand::Stop => self.stop(),
1384 PlayerCommand::Seek(pos) => self.seek(pos),
1385 PlayerCommand::NextTrack => {
1386 let now = std::time::Instant::now();
1388 if now.duration_since(self.last_skip).as_millis() >= 150 {
1389 self.last_skip = now;
1390 self.next_track();
1391 }
1392 }
1393 PlayerCommand::PrevTrack => {
1394 let now = std::time::Instant::now();
1395 if now.duration_since(self.last_skip).as_millis() >= 150 {
1396 self.last_skip = now;
1397 self.prev_track();
1398 }
1399 }
1400 PlayerCommand::AddToPlaylist(items) => {
1401 let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
1402 self.shared_state.add_items(items);
1403 self.push_undo(UndoEntry::Added { ids });
1404 }
1405 PlayerCommand::UpdatePaths(updates) => {
1406 self.shared_state.update_paths(&updates);
1407 if let Some(info) = self.shared_state.track_info()
1408 && let Some((_, new_path)) = updates.iter().find(|(id, _)| *id == info.id)
1409 {
1410 self.shared_state.set_track_info(Some(TrackInfo {
1411 path: new_path.clone(),
1412 ..info
1413 }));
1414 }
1415 }
1416 PlayerCommand::InsertInPlaylist { items, after } => {
1417 let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
1418 self.shared_state.insert_items_after(items, after);
1419 self.push_undo(UndoEntry::Inserted { ids });
1420 }
1421 PlayerCommand::ClearPlaylist => {
1422 self.stop_playback_and_clear_state();
1426 let (items, cursor) = self.shared_state.snapshot_playlist();
1427 self.shared_state.clear_playlist();
1428 self.push_undo(UndoEntry::Replaced { items, cursor });
1429 }
1430 PlayerCommand::ReplacePlaylist { items, start } => {
1431 self.stop_playback_and_clear_state();
1435 let (old_items, cursor) = self.shared_state.snapshot_playlist();
1436 self.shared_state.clear_playlist();
1437 self.push_undo(UndoEntry::Replaced {
1438 items: old_items,
1439 cursor,
1440 });
1441
1442 if items.is_empty() {
1443 return;
1444 }
1445 let start_id = items.get(start).unwrap_or(&items[0]).id;
1446 self.shared_state.add_items(items);
1447 self.play(start_id);
1448 }
1449 PlayerCommand::RemoveFromPlaylist(id) => {
1450 let item = self.shared_state.get_item(id);
1451 let after = self.shared_state.item_before(id);
1452 self.remove_from_playlist(id);
1453 if let Some(item) = item {
1454 self.push_undo(UndoEntry::Removed {
1455 items: vec![(Box::new(item), after)],
1456 });
1457 }
1458 }
1459 PlayerCommand::RemoveFromPlaylistBatch(ids) => {
1460 let items_with_pos = self.snapshot_for_undo(&ids);
1464 let resume_after = match self.shared_state.cursor() {
1465 Some(cursor) if ids.contains(&cursor) => {
1466 Some(self.shared_state.surviving_item_before(cursor, &ids))
1467 }
1468 _ => None,
1469 };
1470
1471 self.shared_state.remove_items(&ids);
1472
1473 if let Some(resume_after) = resume_after {
1474 self.shared_state.set_cursor(resume_after);
1475 self.next_track();
1476 }
1477
1478 if !items_with_pos.is_empty() {
1479 self.push_undo(UndoEntry::Removed {
1480 items: items_with_pos,
1481 });
1482 }
1483 }
1484 PlayerCommand::MoveInPlaylist { id, target, after } => {
1485 let was_after = self.shared_state.item_before(id);
1486 self.shared_state.move_item(id, target, after);
1487 self.push_undo(UndoEntry::Moved { id, was_after });
1488 }
1489 PlayerCommand::MoveItemsInPlaylist { ids, target, after } => {
1490 let entries = self.shared_state.items_before(&ids);
1491 self.shared_state.move_items(&ids, target, after);
1492 self.push_undo(UndoEntry::MovedBatch { entries });
1493 }
1494 PlayerCommand::ReorderPlaylist(order) => {
1495 let entries = self.shared_state.items_before(&order);
1499 self.shared_state.reorder_to(&order);
1500 self.push_undo(UndoEntry::MovedBatch { entries });
1501 }
1502 PlayerCommand::TrackReady(id) => self.track_ready(id),
1503 PlayerCommand::DecodeFinished => self.on_decode_finished(),
1504 PlayerCommand::TrackStreamReady(id) => self.track_stream_ready(id),
1505 PlayerCommand::StreamProbed { id, info, mode } => self.stream_probed(id, *info, mode),
1506 PlayerCommand::TrackFailed(id) => self.track_failed(id),
1507 PlayerCommand::Undo => self.execute_undo(),
1508 PlayerCommand::Redo => self.execute_redo(),
1509 PlayerCommand::BeginUndoBatch => {
1510 self.batch_buffer = Some(Vec::new());
1511 }
1512 PlayerCommand::EndUndoBatch => {
1513 if let Some(entries) = self.batch_buffer.take() {
1514 if entries.len() == 1 {
1515 self.undo_stack.push(entries.into_iter().next().unwrap());
1517 } else if !entries.is_empty() {
1518 self.undo_stack.push(UndoEntry::Batch(entries));
1519 }
1520 }
1521 }
1522 PlayerCommand::SetOutputDevice(name) => self.set_output_device(name),
1523 PlayerCommand::RestartOutput => {
1524 log::info!("restarting audio output");
1525 self.restart_on_current_track();
1526 }
1527 PlayerCommand::ClearOutputDevice => self.clear_output_device(),
1528 }
1529 }
1530
1531 fn apply_entry(&mut self, entry: UndoEntry) -> Option<UndoEntry> {
1533 match entry {
1534 UndoEntry::Added { ids } => {
1535 let items_with_pos = self.snapshot_for_undo(&ids);
1537 self.shared_state.remove_items(&ids);
1538 Some(UndoEntry::Removed {
1539 items: items_with_pos,
1540 })
1541 }
1542 UndoEntry::Removed { items } => {
1543 let mut ids = Vec::with_capacity(items.len());
1545 for (item, after) in items {
1546 ids.push(item.id);
1547 self.shared_state.insert_item_at(*item, after);
1548 }
1549 Some(UndoEntry::Added { ids })
1550 }
1551 UndoEntry::Inserted { ids } => {
1552 let items_with_pos = self.snapshot_for_undo(&ids);
1554 self.shared_state.remove_items(&ids);
1555 Some(UndoEntry::Removed {
1556 items: items_with_pos,
1557 })
1558 }
1559 UndoEntry::Moved { id, was_after } => {
1560 let current_after = self.shared_state.item_before(id);
1561 self.shared_state.move_item_to(id, was_after);
1562 Some(UndoEntry::Moved {
1563 id,
1564 was_after: current_after,
1565 })
1566 }
1567 UndoEntry::MovedBatch { entries } => {
1568 let ids: Vec<QueueItemId> = entries.iter().map(|(id, _)| *id).collect();
1569 let current_positions = self.shared_state.items_before(&ids);
1570 self.shared_state.move_items_to(&entries);
1571 Some(UndoEntry::MovedBatch {
1572 entries: current_positions,
1573 })
1574 }
1575 UndoEntry::Replaced { items, cursor } => {
1576 let (current_items, current_cursor) = self.shared_state.snapshot_playlist();
1577 self.shared_state.restore_playlist(items, cursor);
1578 Some(UndoEntry::Replaced {
1579 items: current_items,
1580 cursor: current_cursor,
1581 })
1582 }
1583 UndoEntry::Batch(entries) => {
1584 let mut inverses = Vec::with_capacity(entries.len());
1586 for entry in entries.into_iter().rev() {
1587 if let Some(inverse) = self.apply_entry(entry) {
1588 inverses.push(inverse);
1589 }
1590 }
1591 inverses.reverse();
1592 Some(UndoEntry::Batch(inverses))
1593 }
1594 }
1595 }
1596
1597 fn reconcile_playback(&mut self) {
1611 let Some(playing) = self.shared_state.track_info().map(|t| t.id) else {
1612 return;
1613 };
1614 if self.shared_state.get_item(playing).is_some() {
1615 return;
1616 }
1617 let resume = (self.shared_state.playback_state() == PlaybackState::Playing)
1622 .then(|| self.shared_state.cursor())
1623 .flatten();
1624 self.stop_playback_and_clear_state();
1625 if let Some(id) = resume {
1626 self.play(id);
1627 }
1628 }
1629
1630 fn execute_undo(&mut self) {
1632 let Some(entry) = self.undo_stack.pop_undo() else {
1633 return;
1634 };
1635 if let Some(inverse) = self.apply_entry(entry) {
1636 self.undo_stack.push_redo(inverse);
1637 }
1638 self.reconcile_playback();
1639 }
1640
1641 fn execute_redo(&mut self) {
1643 let Some(entry) = self.undo_stack.pop_redo() else {
1644 return;
1645 };
1646 if let Some(inverse) = self.apply_entry(entry) {
1647 self.undo_stack.push_undo_keep_redo(inverse);
1648 }
1649 self.reconcile_playback();
1650 }
1651
1652 pub fn run(&mut self) {
1654 use std::time::Duration;
1655
1656 let rx = self.commands.rx.clone();
1657 loop {
1658 match rx.recv_timeout(Duration::from_millis(50)) {
1660 Ok(cmd) => self.process_command(cmd),
1661 Err(crossbeam_channel::RecvTimeoutError::Timeout) => {}
1662 Err(crossbeam_channel::RecvTimeoutError::Disconnected) => break,
1663 }
1664 self.update_playback_state();
1665 }
1666 self.stop();
1667 }
1668
1669 pub fn spawn() -> (
1672 Arc<SharedPlayerState>,
1673 Arc<PlaybackTimeline>,
1674 Arc<VizSnapshot>,
1675 crossbeam_channel::Sender<PlayerCommand>,
1676 ) {
1677 let mut player = Self::new();
1678 player.history = PlayRecorder::spawn();
1679 let state = player.shared_state();
1680 let timeline = player.timeline();
1681 let viz_snapshot = player.viz_snapshot();
1682 let tx = player.command_sender();
1683
1684 thread::Builder::new()
1685 .name("koan-player".into())
1686 .spawn(move || player.run())
1687 .expect("failed to spawn player thread");
1688
1689 (state, timeline, viz_snapshot, tx)
1690 }
1691}
1692
1693#[cfg(test)]
1694mod tests {
1695 #[test]
1696 fn a_download_in_progress_is_known_by_its_own_extension() {
1697 use std::path::Path;
1698 assert_eq!(
1700 lengthless_mode_for(Path::new("/c/t.m4a.part")),
1701 streaming::ProbeMode::LengthlessWholeEnd
1702 );
1703 assert_eq!(
1704 lengthless_mode_for(Path::new("/c/t.OPUS.part")),
1705 streaming::ProbeMode::LengthlessWholeEnd
1706 );
1707 assert_eq!(
1708 lengthless_mode_for(Path::new("/c/t.flac.part")),
1709 streaming::ProbeMode::Lengthless
1710 );
1711 assert_eq!(
1712 media_extension(Path::new("/c/t.m4a.part")).as_deref(),
1713 Some("m4a")
1714 );
1715 assert_eq!(
1716 media_extension(Path::new("/c/t.mp3")).as_deref(),
1717 Some("mp3")
1718 );
1719 }
1720
1721 use super::*;
1722 use state::PlaylistItem;
1723 use std::path::PathBuf;
1724 use std::sync::atomic::AtomicU64;
1725
1726 fn make_item(title: &str) -> PlaylistItem {
1727 PlaylistItem {
1728 playlist_entry_id: None,
1729 id: QueueItemId::new(),
1730 db_id: None,
1731 path: PathBuf::from(format!("/music/{title}.flac")),
1732 title: title.to_string(),
1733 artist: String::new(),
1734 album_artist: String::new(),
1735 album: String::new(),
1736 year: None,
1737 codec: None,
1738 track_number: None,
1739 disc: None,
1740 duration_ms: None,
1741 state: ItemState::Ready,
1742 }
1743 }
1744
1745 fn playlist_ids(player: &Player) -> Vec<QueueItemId> {
1746 let (items, _) = player.shared_state.snapshot_playlist();
1747 items.iter().map(|i| i.id).collect()
1748 }
1749
1750 fn playlist_titles(player: &Player) -> Vec<String> {
1751 let (items, _) = player.shared_state.snapshot_playlist();
1752 items.iter().map(|i| i.title.clone()).collect()
1753 }
1754
1755 fn pending_item(title: &str) -> PlaylistItem {
1756 PlaylistItem {
1757 playlist_entry_id: None,
1758 state: ItemState::Pending,
1759 ..make_item(title)
1760 }
1761 }
1762
1763 fn pretend_playing(player: &mut Player, id: QueueItemId) {
1767 let item = player
1768 .shared_state
1769 .get_item(id)
1770 .expect("item is in the queue");
1771 player.shared_state.set_track_info(Some(TrackInfo {
1772 id,
1773 path: item.path,
1774 codec: String::new(),
1775 sample_rate: 44_100,
1776 bit_depth: None,
1777 bitrate_kbps: None,
1778 channels: 2,
1779 duration_ms: 1_000,
1780 }));
1781 player
1782 .shared_state
1783 .set_playback_state(PlaybackState::Playing);
1784 }
1785
1786 fn playing_id(player: &Player) -> Option<QueueItemId> {
1787 player.shared_state.track_info().map(|t| t.id)
1788 }
1789
1790 fn seed(player: &mut Player, n: usize) -> Vec<QueueItemId> {
1792 let items: Vec<_> = (0..n).map(|i| make_item(&format!("t{i}"))).collect();
1793 let ids = items.iter().map(|i| i.id).collect();
1794 player.process_command(PlayerCommand::AddToPlaylist(items));
1795 ids
1796 }
1797
1798 fn listen(player: &mut Player, from_ms: u64, to_ms: u64) {
1802 let mut at = from_ms;
1803 if let Some(f) = player.in_flight.as_mut() {
1804 f.advance(at); }
1806 while at < to_ms {
1807 at = (at + 50).min(to_ms);
1808 if let Some(f) = player.in_flight.as_mut() {
1809 f.advance(at);
1810 }
1811 }
1812 }
1813
1814 fn start(player: &mut Player, track_id: i64) -> QueueItemId {
1815 let id = QueueItemId::new();
1816 player.on_track_changed(id, 0);
1817 player
1819 .in_flight
1820 .as_mut()
1821 .unwrap()
1822 .track_id_for_test(track_id);
1823 id
1824 }
1825
1826 #[test]
1827 fn a_gapless_transition_closes_the_outgoing_track_and_opens_the_next() {
1828 let mut player = Player::new();
1829 start(&mut player, 11);
1830 listen(&mut player, 0, 200_000);
1831
1832 let b = QueueItemId::new();
1833 player.on_track_changed(b, 0);
1834 let f = player
1835 .in_flight
1836 .as_ref()
1837 .expect("the next track is counting");
1838 assert_eq!(f.item, b);
1839 assert_eq!(f.listened_ms(), 0, "and starts from nothing");
1840 }
1841
1842 #[test]
1843 fn a_track_skipped_seconds_in_is_still_history() {
1844 let mut player = Player::new();
1845 start(&mut player, 7);
1846 listen(&mut player, 0, 2_000);
1847
1848 let event = player
1849 .finish_play()
1850 .expect("putting something on is a thing you did, however briefly");
1851 assert!(matches!(
1852 event,
1853 history::PlayEvent::Finished {
1854 track_id: 7,
1855 listened_ms: 2_000
1856 }
1857 ));
1858 }
1859
1860 #[test]
1861 fn a_track_is_closed_out_once() {
1862 let mut player = Player::new();
1863 start(&mut player, 7);
1864 listen(&mut player, 0, 200_000);
1865
1866 assert!(player.finish_play().is_some());
1867 assert!(player.finish_play().is_none());
1868 }
1869
1870 #[test]
1871 fn seeking_around_a_track_does_not_enter_it_twice() {
1872 let mut player = Player::new();
1873 let id = start(&mut player, 7);
1874 listen(&mut player, 0, 120_000);
1875
1876 player.on_track_changed(id, 30_000);
1878 assert_eq!(
1879 player.in_flight.as_ref().unwrap().listened_ms(),
1880 120_000,
1881 "the seek kept the count rather than restarting it"
1882 );
1883 listen(&mut player, 30_000, 40_000);
1884
1885 let Some(history::PlayEvent::Finished { listened_ms, .. }) = player.finish_play() else {
1886 panic!("still one play");
1887 };
1888 assert_eq!(listened_ms, 130_000);
1889 assert!(player.finish_play().is_none());
1890 }
1891
1892 #[test]
1893 fn a_track_that_is_not_in_the_library_is_not_recorded() {
1894 let mut player = Player::new();
1895 let id = QueueItemId::new();
1896 player.on_track_changed(id, 0);
1897 listen(&mut player, 0, 200_000);
1898 assert!(player.finish_play().is_none());
1899 }
1900
1901 #[test]
1902 fn stopping_closes_out_what_was_heard() {
1903 let mut player = Player::new();
1904 start(&mut player, 7);
1905 listen(&mut player, 0, 150_000);
1906
1907 player.stop_playback_and_clear_state();
1908 assert!(player.in_flight.is_none(), "the stop consumed it");
1909 }
1910
1911 #[test]
1912 fn the_server_hears_each_turn_playback_takes() {
1913 use PlaybackReportState::{Paused, Playing, Stopped};
1914 use history::PlaybackReport;
1915
1916 let dir = tempfile::tempdir().unwrap();
1917 let path = dir.path().join("t.wav");
1918 crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
1919
1920 let mut player = Player::new();
1921 player.backend = Box::new(StuckBackend {
1922 rate: 8_000.0,
1923 asked: Default::default(),
1924 });
1925 let (recorder, events) = PlayRecorder::capture();
1926 player.history = Some(recorder);
1927
1928 let item = PlaylistItem {
1929 db_id: Some(5),
1930 path,
1931 ..make_item("t")
1932 };
1933 let id = item.id;
1934 player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
1935 player.process_command(PlayerCommand::Play(id));
1936 player.process_command(PlayerCommand::Pause);
1937 player.process_command(PlayerCommand::Seek(4_000));
1938 player.process_command(PlayerCommand::Resume);
1939 player.process_command(PlayerCommand::Stop);
1940
1941 let report = |state, position_ms| {
1942 PlayEvent::Playback(PlaybackReport {
1943 track_id: 5,
1944 state,
1945 position_ms,
1946 })
1947 };
1948 assert_eq!(
1949 events.try_iter().collect::<Vec<_>>(),
1950 vec![
1951 PlayEvent::Started {
1952 track_id: 5,
1953 position_ms: 0
1954 },
1955 report(Paused, 0),
1956 report(Paused, 4_000),
1957 report(Playing, 4_000),
1958 report(Stopped, 4_000),
1959 PlayEvent::Finished {
1960 track_id: 5,
1961 listened_ms: 0
1962 },
1963 ]
1964 );
1965 }
1966
1967 #[test]
1968 fn removing_the_playing_track_resumes_at_its_successor() {
1969 let mut player = Player::new();
1970 let ids = seed(&mut player, 5);
1971 player.shared_state.set_cursor(Some(ids[2]));
1972
1973 player.process_command(PlayerCommand::RemoveFromPlaylist(ids[2]));
1974
1975 assert_eq!(
1976 player.shared_state.cursor(),
1977 Some(ids[3]),
1978 "playback must continue at the next track, not restart the queue"
1979 );
1980 assert_eq!(player.playback_starts, 1);
1981 }
1982
1983 #[test]
1984 fn removing_the_first_playing_track_resumes_at_the_new_first() {
1985 let mut player = Player::new();
1986 let ids = seed(&mut player, 3);
1987 player.shared_state.set_cursor(Some(ids[0]));
1988
1989 player.process_command(PlayerCommand::RemoveFromPlaylist(ids[0]));
1990
1991 assert_eq!(player.shared_state.cursor(), Some(ids[1]));
1992 }
1993
1994 #[test]
1995 fn next_track_parks_on_a_track_that_has_not_downloaded_yet() {
1996 let mut player = Player::new();
1997 let playing = make_item("playing");
1998 let waiting = pending_item("waiting");
1999 let later = make_item("later");
2000 let (playing_id, waiting_id) = (playing.id, waiting.id);
2001 player.process_command(PlayerCommand::AddToPlaylist(vec![playing, waiting, later]));
2002 player.shared_state.set_cursor(Some(playing_id));
2003
2004 player.process_command(PlayerCommand::DecodeFinished);
2005
2006 assert_eq!(
2007 player.shared_state.cursor(),
2008 Some(waiting_id),
2009 "the cursor parks on the track being fetched"
2010 );
2011 assert_eq!(
2012 player.playback_starts, 0,
2013 "nothing to play until its bytes land"
2014 );
2015
2016 player
2019 .shared_state
2020 .update_item_state(waiting_id, ItemState::Ready);
2021 player.process_command(PlayerCommand::TrackReady(waiting_id));
2022
2023 assert_eq!(player.playback_starts, 1);
2024 assert_eq!(player.shared_state.cursor(), Some(waiting_id));
2025 }
2026
2027 #[test]
2028 fn a_download_that_cannot_land_moves_the_cursor_on() {
2029 let mut player = Player::new();
2030 let waiting = pending_item("waiting");
2031 let later = make_item("later");
2032 let (waiting_id, later_id) = (waiting.id, later.id);
2033 player.process_command(PlayerCommand::AddToPlaylist(vec![waiting, later]));
2034
2035 player.process_command(PlayerCommand::Play(waiting_id));
2036 assert_eq!(player.playback_starts, 0, "nothing to play yet");
2037
2038 player
2040 .shared_state
2041 .update_item_state(waiting_id, ItemState::Failed("remote unavailable".into()));
2042 player.process_command(PlayerCommand::TrackFailed(waiting_id));
2043
2044 assert_eq!(
2045 player.shared_state.cursor(),
2046 Some(later_id),
2047 "the queue moves past a track that can never load"
2048 );
2049 assert_eq!(player.playback_starts, 1);
2050 }
2051
2052 #[test]
2053 fn a_queue_that_can_never_load_stops_rather_than_waiting() {
2054 let mut player = Player::new();
2055 let first = pending_item("first");
2056 let second = pending_item("second");
2057 let (first_id, second_id) = (first.id, second.id);
2058 player.process_command(PlayerCommand::AddToPlaylist(vec![first, second]));
2059
2060 player.process_command(PlayerCommand::Play(first_id));
2061 for id in [first_id, second_id] {
2062 player
2063 .shared_state
2064 .update_item_state(id, ItemState::Failed("remote unavailable".into()));
2065 player.process_command(PlayerCommand::TrackFailed(id));
2066 }
2067
2068 assert_eq!(player.playback_starts, 0);
2069 assert_eq!(
2070 player.shared_state.playback_state(),
2071 PlaybackState::Stopped,
2072 "a stop the UI can see, not an indefinite wait for TrackReady"
2073 );
2074 }
2075
2076 #[test]
2077 fn a_failure_elsewhere_in_the_queue_leaves_the_cursor_alone() {
2078 let mut player = Player::new();
2079 let waiting = pending_item("waiting");
2080 let other = pending_item("other");
2081 let (waiting_id, other_id) = (waiting.id, other.id);
2082 player.process_command(PlayerCommand::AddToPlaylist(vec![waiting, other]));
2083 player.process_command(PlayerCommand::Play(waiting_id));
2084
2085 player
2086 .shared_state
2087 .update_item_state(other_id, ItemState::Failed("remote unavailable".into()));
2088 player.process_command(PlayerCommand::TrackFailed(other_id));
2089
2090 assert_eq!(
2091 player.shared_state.cursor(),
2092 Some(waiting_id),
2093 "a track still downloading keeps the cursor"
2094 );
2095 }
2096
2097 #[test]
2098 fn batch_delete_containing_the_cursor_restarts_the_engine_once() {
2099 let mut player = Player::new();
2100 let ids = seed(&mut player, 5);
2101 player.shared_state.set_cursor(Some(ids[2]));
2102
2103 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![
2104 ids[1], ids[2], ids[3],
2105 ]));
2106
2107 assert_eq!(playlist_titles(&player), vec!["t0", "t4"]);
2108 assert_eq!(player.shared_state.cursor(), Some(ids[4]));
2109 assert_eq!(
2110 player.playback_starts, 1,
2111 "one resume for the whole selection, not one per deleted track"
2112 );
2113 }
2114
2115 #[test]
2116 fn batch_delete_below_the_cursor_leaves_playback_alone() {
2117 let mut player = Player::new();
2118 let ids = seed(&mut player, 4);
2119 player.shared_state.set_cursor(Some(ids[0]));
2120
2121 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![ids[2], ids[3]]));
2122
2123 assert_eq!(player.shared_state.cursor(), Some(ids[0]));
2124 assert_eq!(player.playback_starts, 0);
2125 }
2126
2127 #[test]
2128 fn undo_of_a_batch_delete_restores_the_original_order() {
2129 let mut player = Player::new();
2133 let items = vec![
2134 make_item("A"),
2135 make_item("B"),
2136 make_item("C"),
2137 make_item("D"),
2138 ];
2139 let (b_id, c_id) = (items[1].id, items[2].id);
2140 player.process_command(PlayerCommand::AddToPlaylist(items));
2141
2142 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![c_id, b_id]));
2143 assert_eq!(playlist_titles(&player), vec!["A", "D"]);
2144
2145 player.process_command(PlayerCommand::Undo);
2146 assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
2147 }
2148
2149 #[test]
2152 fn undo_add_removes_items() {
2153 let mut player = Player::new();
2154 let items = vec![make_item("A"), make_item("B")];
2155 let ids: Vec<_> = items.iter().map(|i| i.id).collect();
2156
2157 player.process_command(PlayerCommand::AddToPlaylist(items));
2158 assert_eq!(playlist_ids(&player), ids);
2159 assert!(player.undo_stack().can_undo());
2160
2161 player.process_command(PlayerCommand::Undo);
2162 assert!(playlist_ids(&player).is_empty());
2163 assert!(player.undo_stack().can_redo());
2164 }
2165
2166 #[test]
2167 fn redo_add_restores_items() {
2168 let mut player = Player::new();
2169 let items = vec![make_item("A"), make_item("B")];
2170
2171 player.process_command(PlayerCommand::AddToPlaylist(items));
2172 player.process_command(PlayerCommand::Undo);
2173 assert!(playlist_ids(&player).is_empty());
2174
2175 player.process_command(PlayerCommand::Redo);
2176 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2177 }
2178
2179 #[test]
2182 fn undo_remove_restores_item_at_position() {
2183 let mut player = Player::new();
2184 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2185 let b_id = items[1].id;
2186
2187 player.process_command(PlayerCommand::AddToPlaylist(items));
2188 player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
2189 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2190
2191 player.process_command(PlayerCommand::Undo);
2192 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2193 }
2194
2195 #[test]
2196 fn undo_remove_first_item() {
2197 let mut player = Player::new();
2198 let items = vec![make_item("A"), make_item("B")];
2199 let a_id = items[0].id;
2200
2201 player.process_command(PlayerCommand::AddToPlaylist(items));
2202 player.process_command(PlayerCommand::RemoveFromPlaylist(a_id));
2203 assert_eq!(playlist_titles(&player), vec!["B"]);
2204
2205 player.process_command(PlayerCommand::Undo);
2206 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2207 }
2208
2209 #[test]
2210 fn undo_batch_remove_restores_all() {
2211 let mut player = Player::new();
2212 let items = vec![
2213 make_item("A"),
2214 make_item("B"),
2215 make_item("C"),
2216 make_item("D"),
2217 ];
2218 let b_id = items[1].id;
2219 let c_id = items[2].id;
2220
2221 player.process_command(PlayerCommand::AddToPlaylist(items));
2222 let version_before = player.shared_state.playlist_version();
2223 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![b_id, c_id]));
2224 assert_eq!(playlist_titles(&player), vec!["A", "D"]);
2225 assert_eq!(
2228 player.shared_state.playlist_version(),
2229 version_before + 1,
2230 "batch removal must bump the playlist version exactly once"
2231 );
2232
2233 player.process_command(PlayerCommand::Undo);
2235 assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
2236 }
2237
2238 #[test]
2239 fn redo_batch_remove() {
2240 let mut player = Player::new();
2241 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2242 let a_id = items[0].id;
2243 let b_id = items[1].id;
2244
2245 player.process_command(PlayerCommand::AddToPlaylist(items));
2246 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![a_id, b_id]));
2247 player.process_command(PlayerCommand::Undo);
2248 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2249
2250 player.process_command(PlayerCommand::Redo);
2251 assert_eq!(playlist_titles(&player), vec!["C"]);
2252 }
2253
2254 #[test]
2255 fn redo_remove() {
2256 let mut player = Player::new();
2257 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2258 let b_id = items[1].id;
2259
2260 player.process_command(PlayerCommand::AddToPlaylist(items));
2261 player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
2262 player.process_command(PlayerCommand::Undo);
2263 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2264
2265 player.process_command(PlayerCommand::Redo);
2266 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2267 }
2268
2269 #[test]
2272 fn undo_insert_removes_inserted_items() {
2273 let mut player = Player::new();
2274 let items = vec![make_item("A"), make_item("C")];
2275 let a_id = items[0].id;
2276
2277 player.process_command(PlayerCommand::AddToPlaylist(items));
2278
2279 let inserted = vec![make_item("B")];
2280 player.process_command(PlayerCommand::InsertInPlaylist {
2281 items: inserted,
2282 after: a_id,
2283 });
2284 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2285
2286 player.process_command(PlayerCommand::Undo);
2287 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2288 }
2289
2290 #[test]
2293 fn undo_move_restores_position() {
2294 let mut player = Player::new();
2295 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2296 let a_id = items[0].id;
2297 let c_id = items[2].id;
2298
2299 player.process_command(PlayerCommand::AddToPlaylist(items));
2300
2301 player.process_command(PlayerCommand::MoveInPlaylist {
2303 id: a_id,
2304 target: c_id,
2305 after: true,
2306 });
2307 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2308
2309 player.process_command(PlayerCommand::Undo);
2310 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2311 }
2312
2313 #[test]
2314 fn redo_move() {
2315 let mut player = Player::new();
2316 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2317 let a_id = items[0].id;
2318 let c_id = items[2].id;
2319
2320 player.process_command(PlayerCommand::AddToPlaylist(items));
2321 player.process_command(PlayerCommand::MoveInPlaylist {
2322 id: a_id,
2323 target: c_id,
2324 after: true,
2325 });
2326 player.process_command(PlayerCommand::Undo);
2327 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2328
2329 player.process_command(PlayerCommand::Redo);
2330 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2331 }
2332
2333 #[test]
2336 fn undo_batch_move() {
2337 let mut player = Player::new();
2338 let items = vec![
2339 make_item("A"),
2340 make_item("B"),
2341 make_item("C"),
2342 make_item("D"),
2343 ];
2344 let a_id = items[0].id;
2345 let b_id = items[1].id;
2346 let d_id = items[3].id;
2347
2348 player.process_command(PlayerCommand::AddToPlaylist(items));
2349
2350 player.process_command(PlayerCommand::MoveItemsInPlaylist {
2352 ids: vec![a_id, b_id],
2353 target: d_id,
2354 after: true,
2355 });
2356 assert_eq!(playlist_titles(&player), vec!["C", "D", "A", "B"]);
2357
2358 player.process_command(PlayerCommand::Undo);
2359 assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
2360 }
2361
2362 #[test]
2365 fn undo_clear_restores_playlist() {
2366 let mut player = Player::new();
2367 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2368
2369 player.process_command(PlayerCommand::AddToPlaylist(items));
2370 player.process_command(PlayerCommand::ClearPlaylist);
2371 assert!(playlist_ids(&player).is_empty());
2372
2373 player.process_command(PlayerCommand::Undo);
2374 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2375 }
2376
2377 #[test]
2382 fn undoing_a_replace_does_not_leave_the_engine_on_an_orphaned_track() {
2383 let mut player = Player::new();
2384 let original = seed(&mut player, 3);
2385 player.shared_state.set_cursor(Some(original[0]));
2386 pretend_playing(&mut player, original[0]);
2387
2388 let replacement = vec![make_item("something else")];
2389 let orphan = replacement[0].id;
2390 player.process_command(PlayerCommand::ReplacePlaylist {
2391 items: replacement,
2392 start: 0,
2393 });
2394 pretend_playing(&mut player, orphan);
2396
2397 player.process_command(PlayerCommand::Undo);
2398
2399 assert_eq!(playlist_ids(&player), original, "the queue comes back");
2400 assert!(
2401 player.shared_state.get_item(orphan).is_none(),
2402 "and the replacement is gone from it"
2403 );
2404 assert!(
2405 playing_id(&player).is_none_or(|id| player.shared_state.get_item(id).is_some()),
2406 "so nothing may still be playing out of it"
2407 );
2408 }
2409
2410 #[test]
2412 fn undoing_an_add_does_not_leave_the_engine_on_a_removed_track() {
2413 let mut player = Player::new();
2414 seed(&mut player, 2);
2415 let added = seed(&mut player, 1);
2416 pretend_playing(&mut player, added[0]);
2417
2418 player.process_command(PlayerCommand::Undo);
2419
2420 assert!(player.shared_state.get_item(added[0]).is_none());
2421 assert!(
2422 playing_id(&player).is_none_or(|id| player.shared_state.get_item(id).is_some()),
2423 "the engine cannot be left on the item the undo removed"
2424 );
2425 }
2426
2427 #[test]
2429 fn undoing_a_move_leaves_playback_alone() {
2430 let mut player = Player::new();
2431 let ids = seed(&mut player, 3);
2432 player.shared_state.set_cursor(Some(ids[0]));
2433 pretend_playing(&mut player, ids[0]);
2434 let starts = player.playback_starts;
2435
2436 player.process_command(PlayerCommand::MoveInPlaylist {
2437 id: ids[2],
2438 target: ids[0],
2439 after: false,
2440 });
2441 player.process_command(PlayerCommand::Undo);
2442
2443 assert_eq!(playlist_ids(&player), ids);
2444 assert_eq!(playing_id(&player), Some(ids[0]), "still on the same track");
2445 assert_eq!(player.playback_starts, starts, "and not restarted");
2446 }
2447
2448 #[test]
2449 fn redo_clear() {
2450 let mut player = Player::new();
2451 let items = vec![make_item("A"), make_item("B")];
2452
2453 player.process_command(PlayerCommand::AddToPlaylist(items));
2454 player.process_command(PlayerCommand::ClearPlaylist);
2455 player.process_command(PlayerCommand::Undo);
2456 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2457
2458 player.process_command(PlayerCommand::Redo);
2459 assert!(playlist_ids(&player).is_empty());
2460 }
2461
2462 #[test]
2465 fn multiple_undos_in_sequence() {
2466 let mut player = Player::new();
2467
2468 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("A")]));
2469 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
2470 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("C")]));
2471 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2472
2473 player.process_command(PlayerCommand::Undo);
2474 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2475
2476 player.process_command(PlayerCommand::Undo);
2477 assert_eq!(playlist_titles(&player), vec!["A"]);
2478
2479 player.process_command(PlayerCommand::Undo);
2480 assert!(playlist_ids(&player).is_empty());
2481 }
2482
2483 #[test]
2484 fn undo_redo_undo_cycle() {
2485 let mut player = Player::new();
2486 let items = vec![make_item("A"), make_item("B")];
2487
2488 player.process_command(PlayerCommand::AddToPlaylist(items));
2489 player.process_command(PlayerCommand::Undo);
2490 assert!(playlist_ids(&player).is_empty());
2491
2492 player.process_command(PlayerCommand::Redo);
2493 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2494
2495 player.process_command(PlayerCommand::Undo);
2496 assert!(playlist_ids(&player).is_empty());
2497 }
2498
2499 #[test]
2500 fn new_action_clears_redo_stack() {
2501 let mut player = Player::new();
2502 let items = vec![make_item("A")];
2503
2504 player.process_command(PlayerCommand::AddToPlaylist(items));
2505 player.process_command(PlayerCommand::Undo);
2506 assert!(player.undo_stack().can_redo());
2507
2508 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
2510 assert!(!player.undo_stack().can_redo());
2511 }
2512
2513 #[test]
2514 fn undo_on_empty_stack_is_noop() {
2515 let mut player = Player::new();
2516 player.process_command(PlayerCommand::Undo);
2517 assert!(playlist_ids(&player).is_empty());
2518 }
2519
2520 #[test]
2521 fn redo_on_empty_stack_is_noop() {
2522 let mut player = Player::new();
2523 player.process_command(PlayerCommand::Redo);
2524 assert!(playlist_ids(&player).is_empty());
2525 }
2526
2527 #[test]
2530 fn playback_commands_not_undoable() {
2531 let mut player = Player::new();
2532 player.process_command(PlayerCommand::Pause);
2533 player.process_command(PlayerCommand::Resume);
2534 player.process_command(PlayerCommand::NextTrack);
2535 player.process_command(PlayerCommand::PrevTrack);
2536 assert!(!player.undo_stack().can_undo());
2537 }
2538
2539 #[test]
2540 fn update_paths_not_undoable() {
2541 let mut player = Player::new();
2542 let items = vec![make_item("A")];
2543 let id = items[0].id;
2544 player.process_command(PlayerCommand::AddToPlaylist(items));
2545
2546 let undo_count = player.undo_stack().undo_len();
2547 player.process_command(PlayerCommand::UpdatePaths(vec![(
2548 id,
2549 PathBuf::from("/new/path.flac"),
2550 )]));
2551 assert_eq!(player.undo_stack().undo_len(), undo_count);
2552 }
2553
2554 #[test]
2557 fn add_remove_undo_undo_produces_original() {
2558 let mut player = Player::new();
2559 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2560 let b_id = items[1].id;
2561 let original_titles = vec!["A", "B", "C"];
2562
2563 player.process_command(PlayerCommand::AddToPlaylist(items));
2564 player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
2565 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2566
2567 player.process_command(PlayerCommand::Undo);
2569 assert_eq!(playlist_titles(&player), original_titles);
2570
2571 player.process_command(PlayerCommand::Undo);
2573 assert!(playlist_ids(&player).is_empty());
2574 }
2575
2576 #[test]
2577 fn interleaved_adds_and_moves_undo() {
2578 let mut player = Player::new();
2579 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2580 let a_id = items[0].id;
2581 let c_id = items[2].id;
2582
2583 player.process_command(PlayerCommand::AddToPlaylist(items));
2584
2585 player.process_command(PlayerCommand::MoveInPlaylist {
2587 id: a_id,
2588 target: c_id,
2589 after: true,
2590 });
2591 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2592
2593 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("D")]));
2595 assert_eq!(playlist_titles(&player), vec!["B", "C", "A", "D"]);
2596
2597 player.process_command(PlayerCommand::Undo);
2599 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2600
2601 player.process_command(PlayerCommand::Undo);
2603 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2604 }
2605
2606 #[test]
2611 fn stop_engine_drops_engine_synchronously() {
2612 use std::sync::atomic::{AtomicBool, Ordering};
2613
2614 struct MockEngine {
2615 dropped: Arc<AtomicBool>,
2616 }
2617 impl AudioEngineHandle for MockEngine {
2618 fn start(&self) -> Result<(), BackendError> {
2619 Ok(())
2620 }
2621 fn stop(&self) -> Result<(), BackendError> {
2622 Ok(())
2623 }
2624 fn is_running(&self) -> bool {
2625 false
2626 }
2627 fn fade_out(&self) {}
2628 fn fade_in(&self) -> Result<(), BackendError> {
2629 Ok(())
2630 }
2631 fn is_silent(&self) -> bool {
2632 false
2633 }
2634 }
2635 impl Drop for MockEngine {
2636 fn drop(&mut self) {
2637 self.dropped.store(true, Ordering::SeqCst);
2638 }
2639 }
2640
2641 let dropped = Arc::new(AtomicBool::new(false));
2642
2643 let stop_flag = Arc::new(AtomicBool::new(false));
2645 let decode_handle = buffer::DecodeHandle::new_for_test(stop_flag);
2646
2647 let mut player = Player::new();
2648 player.active_playback = Some(ActivePlayback {
2649 engine: Box::new(MockEngine {
2650 dropped: dropped.clone(),
2651 }),
2652 decode_handle,
2653 stream: None,
2654 _rate_watch: None,
2655 });
2656
2657 player.stop_engine();
2658
2659 assert!(
2663 dropped.load(Ordering::SeqCst),
2664 "AudioEngine must be dropped synchronously in stop_engine (GitHub #89)"
2665 );
2666 }
2667
2668 #[test]
2669 fn abandoning_a_stream_wakes_a_reader_parked_at_the_write_head() {
2670 let live = LiveStream {
2671 feed: crate::remote::downloads::ByteFeed::new(),
2672 abandoned: Default::default(),
2673 };
2674 let feed = live.feed.clone();
2675 let started = std::time::Instant::now();
2676 let reader = thread::spawn(move || {
2677 feed.wait_past(
2678 0,
2679 std::time::Instant::now() + std::time::Duration::from_secs(30),
2680 )
2681 });
2682 thread::sleep(std::time::Duration::from_millis(50));
2683 live.abandon();
2684 reader.join().unwrap();
2685
2686 assert!(live.abandoned.load(std::sync::atomic::Ordering::Acquire));
2687 assert!(started.elapsed() < std::time::Duration::from_secs(5));
2688 }
2689
2690 struct StuckBackend {
2695 rate: f64,
2696 asked: Arc<std::sync::Mutex<Option<(f64, u32)>>>,
2697 }
2698
2699 struct NullEngine;
2700 impl AudioEngineHandle for NullEngine {
2701 fn start(&self) -> Result<(), BackendError> {
2702 Ok(())
2703 }
2704 fn stop(&self) -> Result<(), BackendError> {
2705 Ok(())
2706 }
2707 fn is_running(&self) -> bool {
2708 false
2709 }
2710 fn fade_out(&self) {}
2711 fn fade_in(&self) -> Result<(), BackendError> {
2712 Ok(())
2713 }
2714 fn is_silent(&self) -> bool {
2715 false
2716 }
2717 }
2718
2719 impl AudioBackend for StuckBackend {
2720 fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
2721 Ok(vec![self.default_device()?])
2722 }
2723 fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
2724 Ok(backend::DeviceInfo {
2725 name: "Stuck DAC".into(),
2726 sample_rates: vec![self.rate],
2727 platform_id: 0,
2728 })
2729 }
2730 fn supported_sample_rates(
2731 &self,
2732 _device: &backend::DeviceInfo,
2733 ) -> Result<Vec<f64>, BackendError> {
2734 Ok(vec![self.rate])
2735 }
2736 fn get_device_sample_rate(
2737 &self,
2738 _device: &backend::DeviceInfo,
2739 ) -> Result<f64, BackendError> {
2740 Ok(self.rate)
2741 }
2742 fn set_device_sample_rate(
2743 &self,
2744 _device: &backend::DeviceInfo,
2745 rate: f64,
2746 ) -> Result<f64, BackendError> {
2747 Err(BackendError::UnsupportedSampleRate(rate))
2748 }
2749 fn create_engine(
2750 &self,
2751 _device: &backend::DeviceInfo,
2752 sample_rate: f64,
2753 channels: u32,
2754 _consumer: rtrb::Consumer<f32>,
2755 _samples_played: Arc<AtomicU64>,
2756 ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
2757 *self.asked.lock().unwrap() = Some((sample_rate, channels));
2758 Ok(Box::new(NullEngine))
2759 }
2760 }
2761
2762 fn engine_format_for(source_rate: u32, channels: u16, device_rate: f64) -> (f64, u32) {
2763 let asked = Arc::new(std::sync::Mutex::new(None));
2764 let mut player = Player::new();
2765 player.backend = Box::new(StuckBackend {
2766 rate: device_rate,
2767 asked: asked.clone(),
2768 });
2769
2770 let info = buffer::StreamInfo {
2771 codec: "MP3".into(),
2772 sample_rate: source_rate,
2773 channels,
2774 bit_depth: Some(16),
2775 bitrate_kbps: None,
2776 duration_ms: 1000,
2777 };
2778 let (_producer, consumer) = rtrb::RingBuffer::new(16);
2779 player
2780 .create_engine_for(&info, consumer)
2781 .expect("engine creation should succeed");
2782 let asked = *asked.lock().unwrap();
2783 asked.expect("engine was never created")
2784 }
2785
2786 #[test]
2787 fn engine_uses_source_rate_when_device_refuses_switch() {
2788 assert_eq!(engine_format_for(22050, 2, 48000.0), (22050.0, 2));
2791 assert_eq!(engine_format_for(32000, 2, 44100.0), (32000.0, 2));
2792 }
2793
2794 #[test]
2795 fn engine_uses_source_channel_count() {
2796 assert_eq!(engine_format_for(44100, 1, 44100.0), (44100.0, 1));
2797 }
2798
2799 fn settled_rate_for(source_rate: u32, device_rate: f64) -> Option<u32> {
2801 let mut player = Player::new();
2802 player.backend = Box::new(StuckBackend {
2803 rate: device_rate,
2804 asked: Arc::new(std::sync::Mutex::new(None)),
2805 });
2806 let state = player.shared_state.clone();
2807
2808 let info = buffer::StreamInfo {
2809 codec: "MP3".into(),
2810 sample_rate: source_rate,
2811 channels: 2,
2812 bit_depth: Some(16),
2813 bitrate_kbps: None,
2814 duration_ms: 1000,
2815 };
2816 let (_producer, consumer) = rtrb::RingBuffer::new(16);
2817 player
2818 .create_engine_for(&info, consumer)
2819 .expect("engine creation should succeed");
2820 state.output_sample_rate()
2821 }
2822
2823 #[test]
2824 fn settled_device_rate_reaches_the_shared_state() {
2825 assert_eq!(settled_rate_for(22050, 48000.0), Some(48000));
2829 assert_eq!(settled_rate_for(44100, 44100.0), Some(44100));
2831 }
2832
2833 struct SlowBackend {
2835 observed: Arc<std::sync::Mutex<Vec<Option<u32>>>>,
2836 state: Arc<SharedPlayerState>,
2837 }
2838
2839 impl AudioBackend for SlowBackend {
2840 fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
2841 Ok(vec![self.default_device()?])
2842 }
2843 fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
2844 Ok(backend::DeviceInfo {
2845 name: "Slow DAC".into(),
2846 sample_rates: vec![44100.0, 48000.0],
2847 platform_id: 0,
2848 })
2849 }
2850 fn supported_sample_rates(
2851 &self,
2852 _device: &backend::DeviceInfo,
2853 ) -> Result<Vec<f64>, BackendError> {
2854 Ok(vec![44100.0, 48000.0])
2855 }
2856 fn get_device_sample_rate(
2857 &self,
2858 _device: &backend::DeviceInfo,
2859 ) -> Result<f64, BackendError> {
2860 Ok(48000.0)
2861 }
2862 fn set_device_sample_rate(
2863 &self,
2864 _device: &backend::DeviceInfo,
2865 rate: f64,
2866 ) -> Result<f64, BackendError> {
2867 self.observed
2869 .lock()
2870 .unwrap()
2871 .push(self.state.output_sample_rate());
2872 Ok(rate)
2873 }
2874 fn create_engine(
2875 &self,
2876 _device: &backend::DeviceInfo,
2877 _sample_rate: f64,
2878 _channels: u32,
2879 _consumer: rtrb::Consumer<f32>,
2880 _samples_played: Arc<AtomicU64>,
2881 ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
2882 Ok(Box::new(NullEngine))
2883 }
2884 }
2885
2886 #[test]
2887 fn the_previous_rate_is_not_published_while_the_device_reclocks() {
2888 let mut player = Player::new();
2894 let state = player.shared_state.clone();
2895 state.set_output_sample_rate(48000);
2896
2897 let observed = Arc::new(std::sync::Mutex::new(Vec::new()));
2898 player.backend = Box::new(SlowBackend {
2899 observed: observed.clone(),
2900 state: state.clone(),
2901 });
2902
2903 let info = buffer::StreamInfo {
2904 codec: "FLAC".into(),
2905 sample_rate: 44100,
2906 channels: 2,
2907 bit_depth: Some(16),
2908 bitrate_kbps: None,
2909 duration_ms: 1000,
2910 };
2911 let (_producer, consumer) = rtrb::RingBuffer::new(16);
2912 player
2913 .create_engine_for(&info, consumer)
2914 .expect("engine creation should succeed");
2915
2916 assert_eq!(
2917 *observed.lock().unwrap(),
2918 vec![None],
2919 "mid-switch the output rate must read as unknown, not as the last track's"
2920 );
2921 assert_eq!(state.output_sample_rate(), Some(44100));
2922 }
2923
2924 struct WatchedBackend {
2926 inner: StuckBackend,
2927 #[allow(clippy::type_complexity)]
2928 captured: Arc<std::sync::Mutex<Option<Box<dyn Fn(f64) + Send + Sync>>>>,
2929 }
2930
2931 struct NullWatch;
2932 impl backend::SampleRateWatch for NullWatch {}
2933
2934 impl AudioBackend for WatchedBackend {
2935 fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
2936 self.inner.list_devices()
2937 }
2938 fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
2939 self.inner.default_device()
2940 }
2941 fn supported_sample_rates(
2942 &self,
2943 device: &backend::DeviceInfo,
2944 ) -> Result<Vec<f64>, BackendError> {
2945 self.inner.supported_sample_rates(device)
2946 }
2947 fn get_device_sample_rate(
2948 &self,
2949 device: &backend::DeviceInfo,
2950 ) -> Result<f64, BackendError> {
2951 self.inner.get_device_sample_rate(device)
2952 }
2953 fn set_device_sample_rate(
2954 &self,
2955 device: &backend::DeviceInfo,
2956 rate: f64,
2957 ) -> Result<f64, BackendError> {
2958 self.inner.set_device_sample_rate(device, rate)
2959 }
2960 fn watch_device_sample_rate(
2961 &self,
2962 _device: &backend::DeviceInfo,
2963 on_change: Box<dyn Fn(f64) + Send + Sync>,
2964 ) -> Option<Box<dyn backend::SampleRateWatch>> {
2965 *self.captured.lock().unwrap() = Some(on_change);
2966 Some(Box::new(NullWatch))
2967 }
2968 fn create_engine(
2969 &self,
2970 device: &backend::DeviceInfo,
2971 sample_rate: f64,
2972 channels: u32,
2973 consumer: rtrb::Consumer<f32>,
2974 samples_played: Arc<AtomicU64>,
2975 ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
2976 self.inner
2977 .create_engine(device, sample_rate, channels, consumer, samples_played)
2978 }
2979 }
2980
2981 #[test]
2982 fn external_rate_change_reaches_the_shared_state() {
2983 let captured = Arc::new(std::sync::Mutex::new(None));
2987 let mut player = Player::new();
2988 player.backend = Box::new(WatchedBackend {
2989 inner: StuckBackend {
2990 rate: 44100.0,
2991 asked: Arc::new(std::sync::Mutex::new(None)),
2992 },
2993 captured: captured.clone(),
2994 });
2995 let state = player.shared_state.clone();
2996
2997 let info = buffer::StreamInfo {
2998 codec: "FLAC".into(),
2999 sample_rate: 44100,
3000 channels: 2,
3001 bit_depth: Some(16),
3002 bitrate_kbps: None,
3003 duration_ms: 1000,
3004 };
3005 let (_producer, consumer) = rtrb::RingBuffer::new(16);
3006 player
3007 .create_engine_for(&info, consumer)
3008 .expect("engine creation should succeed");
3009 assert_eq!(state.output_sample_rate(), Some(44100));
3010
3011 let on_change = captured.lock().unwrap().take().expect("watch registered");
3012 on_change(48000.0);
3013 assert_eq!(state.output_sample_rate(), Some(48000));
3014 }
3015}