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