1pub mod commands;
2pub mod history;
3mod renderer;
4pub mod state;
5pub mod undo;
6
7use std::path::{Path, PathBuf};
8use std::sync::Arc;
9use std::thread;
10
11use thiserror::Error;
12
13use crate::audio::{
14 analyzer::VizAnalyzer,
15 backend::{self, AudioBackend, AudioEngineHandle, BackendError, SampleRateWatch},
16 buffer, streaming,
17 viz::{VizBuffer, VizSnapshot},
18};
19use crate::remote::client::PlaybackReportState;
20use buffer::PlaybackTimeline;
21use commands::{CommandChannel, PlayerCommand};
22use history::{InFlight, PlayEvent, PlayRecorder, PlaybackReport};
23use state::{
24 ItemState, PlayMode, PlaybackSource, PlaybackState, QueueItemId, Repeat, SharedPlayerState,
25 TrackInfo,
26};
27use undo::{UndoEntry, UndoStack};
28
29pub(crate) const RING_BUFFER_SIZE: usize = 192_000 * 2;
31
32const SEEK_END_GUARD_MS: u64 = 500;
35
36const BOUNDARY_SLACK: std::time::Duration = std::time::Duration::from_millis(5);
39
40const FADE_CHECK: std::time::Duration = std::time::Duration::from_millis(50);
43
44#[derive(Debug, Error)]
45pub enum PlayerError {
46 #[error("backend error: {0}")]
47 Backend(#[from] BackendError),
48 #[error("decode error: {0}")]
49 Decode(#[from] buffer::DecodeError),
50 #[error("renderer: {0}")]
51 Renderer(String),
52 #[error("{0}")]
55 Unplayable(String),
56}
57
58#[derive(Clone)]
61struct StreamSource {
62 path: PathBuf,
63 bytes_written: Arc<crate::remote::downloads::ByteFeed>,
64 total: u64,
65 mode: streaming::ProbeMode,
66}
67
68fn media_extension(path: &Path) -> Option<String> {
71 crate::remote::download::strip_part_suffix(path)
72 .extension()
73 .and_then(|e| e.to_str())
74 .map(str::to_ascii_lowercase)
75}
76
77fn hint_for(path: &Path) -> symphonia::core::formats::probe::Hint {
79 let mut hint = symphonia::core::formats::probe::Hint::new();
80 if let Some(ext) = media_extension(path) {
81 hint.with_extension(&ext);
82 }
83 hint
84}
85
86fn lengthless_mode_for(path: &Path) -> streaming::ProbeMode {
96 match media_extension(path).as_deref() {
97 Some("ogg" | "oga" | "opus" | "spx" | "m4a" | "m4b" | "mp4" | "mov") => {
98 streaming::ProbeMode::LengthlessWholeEnd
99 }
100 _ => streaming::ProbeMode::Lengthless,
101 }
102}
103
104pub struct Player {
106 shared_state: Arc<SharedPlayerState>,
107 commands: CommandChannel,
108 transport: Transport,
110 timeline: Arc<PlaybackTimeline>,
111 viz_buffer: Arc<VizBuffer>,
112 viz_snapshot: Arc<VizSnapshot>,
113 _viz_analyzer: VizAnalyzer,
115 undo_stack: UndoStack,
116 batch_buffer: Option<Vec<UndoEntry>>,
119 output_device_name: Option<String>,
121 backend: Box<dyn AudioBackend>,
123 stream_mode: streaming::ProbeMode,
126 history: Option<PlayRecorder>,
129 downloads: Option<crate::remote::queue::DownloadQueue>,
133 in_flight: Option<InFlight>,
135 lead_in_ends: Option<std::time::Instant>,
137 session: u64,
140 silence_waiters: Vec<crossbeam_channel::Sender<u64>>,
142 dsp: Option<DspCache>,
145 renderer: Option<renderer::RendererLink>,
148 mode: PlayMode,
151 #[cfg(test)]
154 playback_starts: usize,
155 #[cfg(test)]
159 dsp_override: Option<Arc<crate::audio::dsp::Setup>>,
160 #[cfg(test)]
161 renderer_dsp_override: Option<Arc<crate::audio::dsp::Setup>>,
162 resume_renderer: bool,
165}
166
167struct DspCache {
168 config: Arc<crate::config::Config>,
169 device: String,
170 setup: Option<Arc<crate::audio::dsp::Setup>>,
171}
172
173#[derive(Clone, Copy, PartialEq, Eq)]
176enum Run {
177 Playing,
178 Paused,
179}
180
181enum Transport {
185 Idle,
187 Waiting(Waiting),
190 Loaded(Session),
192}
193
194#[derive(Clone, Copy)]
199struct Waiting {
200 id: QueueItemId,
201 position_ms: u64,
202 start: Run,
203}
204
205impl Waiting {
206 fn may_stream(&self, id: QueueItemId) -> bool {
207 self.id == id && self.position_ms == 0
208 }
209}
210
211struct Session {
214 track: TrackInfo,
217 run: Run,
218 lookahead: Arc<parking_lot::Mutex<Vec<state::Lookahead>>>,
221 output: Output,
222}
223
224enum Output {
226 Local(Local),
229 Renderer(Box<renderer::Play>),
233}
234
235struct Local {
236 engine: Box<dyn AudioEngineHandle>,
237 decode_handle: buffer::DecodeHandle,
238 stream: Option<LiveStream>,
240 _rate_watch: Option<Box<dyn SampleRateWatch>>,
243 dsp: Option<crate::audio::dsp::DspStatus>,
245}
246
247impl Session {
248 fn engine(&self) -> Option<&dyn AudioEngineHandle> {
250 match &self.output {
251 Output::Local(local) => Some(local.engine.as_ref()),
252 Output::Renderer(_) => None,
253 }
254 }
255}
256
257enum Source {
259 File(PathBuf),
260 Stream(StreamSource),
261}
262
263impl Source {
264 fn path(&self) -> &Path {
265 match self {
266 Source::File(path) => path,
267 Source::Stream(source) => &source.path,
268 }
269 }
270}
271
272struct LiveStream {
276 feed: Arc<crate::remote::downloads::ByteFeed>,
277 abandoned: Arc<std::sync::atomic::AtomicBool>,
278}
279
280impl LiveStream {
281 fn abandon(&self) {
282 self.abandoned
283 .store(true, std::sync::atomic::Ordering::Release);
284 self.feed.done();
285 }
286}
287
288impl Default for Player {
289 fn default() -> Self {
290 Self::new()
291 }
292}
293
294impl Player {
295 pub fn new() -> Self {
296 let viz_buffer = VizBuffer::new();
297 let viz_snapshot = VizSnapshot::new();
298 let timeline = PlaybackTimeline::new();
299 let cfg = crate::config::Config::cached();
300 let viz_analyzer = VizAnalyzer::spawn_with_snapshot(
301 Arc::clone(&viz_buffer),
302 &cfg.visualizer,
303 Arc::clone(&viz_snapshot),
304 timeline.samples_played_counter(),
305 );
306
307 let shared_state = SharedPlayerState::new();
308 shared_state.attach_timeline(timeline.clone());
309 let commands = CommandChannel::new();
310 let tx = commands.tx.clone();
311 timeline.on_queued(move || {
312 let _ = tx.try_send(PlayerCommand::TrackQueued);
313 });
314
315 Self {
316 shared_state,
317 commands,
318 transport: Transport::Idle,
319 lead_in_ends: None,
320 dsp: None,
321 session: 0,
322 silence_waiters: Vec::new(),
323 renderer: None,
324 mode: PlayMode::default(),
325 timeline,
326 viz_buffer,
327 viz_snapshot,
328 _viz_analyzer: viz_analyzer,
329 undo_stack: UndoStack::new(),
330 batch_buffer: None,
331 output_device_name: cfg.playback.output_device.clone(),
332 backend: crate::audio::platform_backend(),
333 stream_mode: streaming::ProbeMode::Full,
334 history: None,
335 downloads: None,
336 in_flight: None,
337 #[cfg(test)]
338 playback_starts: 0,
339 #[cfg(test)]
340 dsp_override: None,
341 #[cfg(test)]
342 renderer_dsp_override: None,
343 resume_renderer: false,
344 }
345 }
346
347 pub fn shared_state(&self) -> Arc<SharedPlayerState> {
349 self.shared_state.clone()
350 }
351
352 pub fn timeline(&self) -> Arc<PlaybackTimeline> {
354 self.timeline.clone()
355 }
356
357 pub fn viz_buffer(&self) -> Arc<VizBuffer> {
359 self.viz_buffer.clone()
360 }
361
362 pub fn viz_snapshot(&self) -> Arc<VizSnapshot> {
365 self.viz_snapshot.clone()
366 }
367
368 pub fn undo_stack(&self) -> &UndoStack {
370 &self.undo_stack
371 }
372
373 #[allow(clippy::type_complexity)]
381 fn create_engine_for(
382 &mut self,
383 info: &buffer::StreamInfo,
384 consumer: rtrb::Consumer<f32>,
385 ) -> Result<(Box<dyn AudioEngineHandle>, Option<Box<dyn SampleRateWatch>>), PlayerError> {
386 let device = self.resolve_device()?;
387 let device_rate = self.backend.get_device_sample_rate(&device)?;
388 let dsp = match &self.dsp {
391 Some(cache) if cache.device == device.name => cache.setup.clone(),
392 _ => self.dsp_for(&device.name),
393 };
394 let source_rate =
397 dsp.as_ref()
398 .map_or(info.sample_rate, |d| d.output_rate(info.sample_rate)) as f64;
399
400 self.shared_state.clear_output_sample_rate();
404
405 let settled = if (device_rate - source_rate).abs() > 0.1 {
406 log::info!(
407 "switching device sample rate: {}Hz → {}Hz",
408 device_rate,
409 source_rate
410 );
411 match self.backend.set_device_sample_rate(&device, source_rate) {
412 Ok(rate) => rate,
413 Err(e) => {
414 log::warn!("failed to set device sample rate: {}", e);
415 device_rate
416 }
417 }
418 } else {
419 device_rate
420 };
421
422 if (settled - source_rate).abs() > 0.1 {
423 log::warn!(
424 "device stayed at {}Hz (wanted {}Hz) — output is resampled, not bit-perfect",
425 settled,
426 source_rate
427 );
428 }
429
430 self.shared_state
433 .set_output_sample_rate(settled.round() as u32);
434
435 let watch_state = self.shared_state.clone();
439 let watch_name = device.name.clone();
440 let rate_watch = self.backend.watch_device_sample_rate(
441 &device,
442 Box::new(move |rate| {
443 log::info!("device sample rate changed externally: {rate}Hz on '{watch_name}'");
444 watch_state.set_output_sample_rate(rate.round() as u32);
445 }),
446 );
447
448 let engine = self.backend.create_engine(
449 &device,
450 source_rate,
451 info.channels as u32,
452 consumer,
453 self.timeline.samples_played_counter(),
454 )?;
455 let now = std::time::Instant::now();
459 let lead_in = if device_rate > 0.0 && (settled - device_rate).abs() > 0.1 {
460 std::time::Duration::from_millis(
461 crate::config::Config::cached()
462 .playback
463 .rate_switch_lead_in_ms as u64,
464 )
465 } else {
466 self.lead_in_ends
467 .map(|end| end.saturating_duration_since(now))
468 .unwrap_or_default()
469 };
470 self.lead_in_ends = None;
471 if !lead_in.is_zero() {
472 engine.lead_in((settled * lead_in.as_secs_f64()) as u64);
473 self.lead_in_ends = Some(now + lead_in);
474 }
475
476 Ok((engine, rate_watch))
477 }
478
479 fn dsp_for(&mut self, device: &str) -> Option<Arc<crate::audio::dsp::Setup>> {
481 let config = crate::config::Config::cached();
482 #[cfg(test)]
483 if let Some(setup) = if self
484 .renderer
485 .as_ref()
486 .is_some_and(|l| l.device_name() == device)
487 {
488 self.renderer_dsp_override.clone()
489 } else {
490 self.dsp_override.clone()
491 } {
492 self.dsp = Some(DspCache {
493 config,
494 device: device.to_string(),
495 setup: Some(setup.clone()),
496 });
497 return Some(setup);
498 }
499 if let Some(cache) = &self.dsp
500 && Arc::ptr_eq(&cache.config, &config)
501 && cache.device == device
502 {
503 return cache.setup.clone();
504 }
505 let setup = config.dsp.profile_for(device).and_then(|profile| {
506 crate::audio::dsp::Setup::load(profile, &crate::config::config_dir())
507 .inspect_err(|e| {
508 log::error!(
509 "dsp: profile '{}' not loaded, playing without it: {e}",
510 profile.name
511 )
512 })
513 .ok()
514 .flatten()
515 .map(Arc::new)
516 });
517 self.dsp = Some(DspCache {
518 config,
519 device: device.to_string(),
520 setup: setup.clone(),
521 });
522 setup
523 }
524
525 fn reload_dsp(&mut self) {
531 let device = match &self.renderer {
532 Some(link) => Ok(link.device_name().to_string()),
533 None => self.resolve_device().map(|d| d.name),
534 };
535 let Ok(device) = device else {
536 self.dsp = None;
537 return;
538 };
539 let was = self
540 .dsp
541 .take()
542 .filter(|c| c.device == device)
543 .map(|c| c.setup);
544 let now = self.dsp_for(&device);
545 let changed = match was {
546 Some(was) => was != now,
547 None => now.is_some(),
548 };
549 if changed {
550 self.restart_on_current_track();
551 }
552 }
553
554 fn processing(&mut self) -> buffer::Processing {
556 let cfg = crate::config::Config::cached();
557 let dsp = match self.resolve_device() {
558 Ok(device) => self.dsp_for(&device.name),
559 Err(_) => None,
560 };
561 buffer::Processing {
562 rg_mode: cfg.playback.replaygain,
563 pre_amp_db: cfg.playback.pre_amp_db,
564 dsp,
565 }
566 }
567
568 fn resolve_device(&self) -> Result<backend::DeviceInfo, PlayerError> {
571 if let Some(ref name) = self.output_device_name {
572 match self.backend.list_devices() {
573 Ok(devices) => {
574 if let Some(dev) = devices.into_iter().find(|d| d.name == *name) {
575 return Ok(dev);
576 }
577 log::warn!(
578 "configured output device '{}' not found, falling back to default",
579 name,
580 );
581 }
582 Err(e) => {
583 log::warn!("failed to list devices while resolving '{}': {}", name, e);
584 }
585 }
586 }
587 Ok(self.backend.default_device()?)
588 }
589
590 pub fn set_output_device(&mut self, name: String) {
593 log::info!("switching output device to: {}", name);
594 self.output_device_name = Some(name.clone());
595
596 if let Err(e) = crate::config::Config::persist(|cfg| {
597 cfg.playback.output_device = Some(name);
598 cfg.playback.renderer = None;
599 cfg.playback.renderer_name = None;
600 }) {
601 log::error!("failed to save output device config: {}", e);
602 }
603
604 if self.renderer.is_some() {
605 self.use_renderer(None);
606 } else {
607 self.restart_on_current_track();
608 }
609 }
610
611 pub fn clear_output_device(&mut self) {
613 log::info!("reverting to system default output device");
614 self.output_device_name = None;
615
616 if let Err(e) = crate::config::Config::persist(|cfg| {
617 cfg.playback.output_device = None;
618 cfg.playback.renderer = None;
619 cfg.playback.renderer_name = None;
620 }) {
621 log::error!("failed to save output device config: {}", e);
622 }
623
624 if self.renderer.is_some() {
625 self.use_renderer(None);
626 } else {
627 self.restart_on_current_track();
628 }
629 }
630
631 fn restart_on_current_track(&mut self) {
634 let position_ms = self.shared_state.position_ms();
635 if let Err(e) = self.restart_current(position_ms) {
636 log::error!("failed to restart playback on device switch: {}", e);
637 }
638 }
639
640 pub fn output_device_name(&self) -> Option<&str> {
642 self.output_device_name.as_deref()
643 }
644
645 pub fn command_sender(&self) -> crossbeam_channel::Sender<PlayerCommand> {
647 self.commands.tx.clone()
648 }
649
650 fn session(&self) -> Option<&Session> {
651 match &self.transport {
652 Transport::Loaded(session) => Some(session),
653 _ => None,
654 }
655 }
656
657 fn waiting(&self) -> Option<Waiting> {
658 match self.transport {
659 Transport::Waiting(waiting) => Some(waiting),
660 _ => None,
661 }
662 }
663
664 fn may_stream(&self, waiting: &Waiting, id: QueueItemId) -> bool {
667 waiting.may_stream(id) && self.renderer.is_none()
668 }
669
670 fn forget_waiting(&mut self) {
671 if matches!(self.transport, Transport::Waiting(_)) {
672 self.transport = Transport::Idle;
673 }
674 }
675
676 fn publish(&self) {
680 let state = &self.shared_state;
681 let (playback, waiting, track) = match &self.transport {
682 Transport::Idle => (PlaybackState::Stopped, false, None),
683 Transport::Waiting(waiting) => (
684 match waiting.start {
685 Run::Playing => PlaybackState::Stopped,
686 Run::Paused => PlaybackState::Paused,
687 },
688 true,
689 None,
690 ),
691 Transport::Loaded(session) => (
692 match session.run {
693 Run::Playing => PlaybackState::Playing,
694 Run::Paused => PlaybackState::Paused,
695 },
696 false,
697 Some(&session.track),
698 ),
699 };
700 if state.track_info().as_ref() != track {
701 state.set_track_info(track.cloned());
702 }
703 state.set_transport(playback, waiting);
704 let dsp = match &self.transport {
707 Transport::Loaded(Session {
708 output: Output::Local(local),
709 ..
710 }) => local.dsp.clone(),
711 Transport::Loaded(Session {
712 output: Output::Renderer(play),
713 ..
714 }) => play.dsp(),
715 _ => None,
716 };
717 if state.dsp() != dsp {
718 state.set_dsp(dsp);
719 }
720 state.set_play_mode(self.mode);
721 }
722
723 fn intent(&self) -> Option<Run> {
728 match &self.transport {
729 Transport::Idle => None,
730 Transport::Waiting(waiting) => Some(waiting.start),
731 Transport::Loaded(session) => Some(session.run),
732 }
733 }
734
735 fn carry_on(&mut self, to: Option<QueueItemId>, intent: Option<Run>) {
738 match (to, intent) {
739 (Some(id), Some(start)) => self.cue(id, 0, start),
740 (Some(id), None) => {
741 self.stop_playback_and_clear_state();
742 self.shared_state.set_cursor(Some(id));
743 }
744 (None, _) => self.stop_playback_and_clear_state(),
745 }
746 }
747
748 pub fn play(&mut self, id: QueueItemId) {
751 if self.in_flight.as_ref().is_some_and(|f| f.item == id) {
754 self.finish_play();
755 }
756 self.forget_waiting();
757 self.shared_state.set_cursor(Some(id));
758
759 match self.shared_state.item_playback_source(id) {
760 Some(PlaybackSource::Ready(path)) => {
761 if let Err(e) = self.start_playback(id, &path, 0, Run::Playing) {
762 log::error!("play failed: {}", e);
763 }
764 }
765 Some(PlaybackSource::Streaming {
766 path,
767 bytes_written,
768 total,
769 }) => {
770 self.park(id, 0, Run::Playing);
775 if self.renderer.is_none() {
776 self.probe_stream_for_playback(id, &path, bytes_written, total);
777 }
778 }
779 None => {
780 self.park(id, 0, Run::Playing);
781 log::info!("play: item {:?} not ready, waiting for TrackReady", id);
782 }
783 }
784 }
785
786 fn cue(&mut self, id: QueueItemId, position_ms: u64, start: Run) {
791 self.forget_waiting();
792 self.shared_state.set_cursor(Some(id));
793 let Some(PlaybackSource::Ready(path)) = self.shared_state.item_playback_source(id) else {
794 if position_ms == 0 && start == Run::Playing {
795 self.play(id);
796 return;
797 }
798 self.park(id, position_ms, start);
799 log::info!("cue: {id:?} not on disk yet, opening at {position_ms}ms once it is");
800 return;
801 };
802 if let Err(e) = self.start_playback(id, &path, position_ms, start) {
803 log::error!("cue failed: {}", e);
804 return;
805 }
806 self.report(match start {
807 Run::Playing => PlaybackReportState::Playing,
808 Run::Paused => PlaybackReportState::Paused,
809 });
810 }
811
812 fn park(&mut self, id: QueueItemId, position_ms: u64, start: Run) {
816 self.report(PlaybackReportState::Stopped);
817 self.finish_play();
818 self.stop_engine();
819 self.timeline.reset();
820 self.shared_state.set_position_ms(position_ms);
821 self.transport = Transport::Waiting(Waiting {
822 id,
823 position_ms,
824 start,
825 });
826 }
827
828 fn start_playback(
829 &mut self,
830 id: QueueItemId,
831 path: &Path,
832 seek_ms: u64,
833 start: Run,
834 ) -> Result<(), PlayerError> {
835 self.open_session(id, Source::File(path.to_path_buf()), None, seek_ms, start)
836 }
837
838 fn open_session(
846 &mut self,
847 id: QueueItemId,
848 source: Source,
849 info: Option<buffer::StreamInfo>,
850 seek_ms: u64,
851 start: Run,
852 ) -> Result<(), PlayerError> {
853 #[cfg(test)]
854 {
855 self.playback_starts += 1;
856 }
857 self.forget_waiting();
858 let mut result = self.try_open_session(id, source, info, seek_ms, start);
859 let mut skipped = id;
865 while matches!(result, Err(PlayerError::Unplayable(_)))
866 && self.shared_state.is_cursor(skipped)
867 {
868 let Some(next) = self.shared_state.advance_cursor_loadable() else {
869 log::info!("upnp: nothing further the renderer can play");
870 self.stop_playback_and_clear_state();
871 return Ok(());
872 };
873 match self.shared_state.item_playback_source(next) {
874 Some(PlaybackSource::Ready(path)) => {
875 skipped = next;
876 result = self.try_open_session(next, Source::File(path), None, 0, start);
877 }
878 _ => {
880 self.cue(next, 0, start);
881 return Ok(());
882 }
883 }
884 }
885 if result.is_err() {
886 self.stop_playback_and_clear_state();
887 }
888 self.wake_analyzer();
889 result
890 }
891
892 fn try_open_session(
893 &mut self,
894 id: QueueItemId,
895 source: Source,
896 info: Option<buffer::StreamInfo>,
897 seek_ms: u64,
898 start: Run,
899 ) -> Result<(), PlayerError> {
900 if self.renderer.is_some() {
901 return self.try_open_on_renderer(id, source, info, seek_ms, start);
902 }
903 self.stop_engine();
904 let info = match info {
905 Some(info) => info,
906 None => buffer::probe_file(source.path())?,
907 };
908 let path = source.path().to_path_buf();
909 let streaming = matches!(source, Source::Stream(_));
910
911 let (first, stream) = match source {
912 Source::File(path) => (buffer::SourceEntry::from_file(id, path), None),
913 Source::Stream(source) => {
914 self.stream_mode = source.mode;
916 let live = LiveStream {
917 feed: source.bytes_written.clone(),
918 abandoned: Default::default(),
919 };
920 let status = {
921 let downloading = self.stream_status_fn(id);
922 let abandoned = live.abandoned.clone();
923 Arc::new(move || {
924 if abandoned.load(std::sync::atomic::Ordering::Acquire) {
925 streaming::StreamStatus::Failed
926 } else {
927 downloading()
928 }
929 }) as Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync>
930 };
931 let entry = buffer::SourceEntry {
932 id,
933 path: source.path.clone(),
934 hint: hint_for(&source.path),
935 make_mss: Box::new(move || {
936 let partial = streaming::PartialFileSource::open(
937 &source.path,
938 source.bytes_written.clone(),
939 source.total,
940 status.clone(),
941 source.mode,
942 )?;
943 Ok(symphonia::core::io::MediaSourceStream::new(
944 Box::new(partial),
945 Default::default(),
946 ))
947 }),
948 };
949 (entry, Some(live))
950 }
951 };
952
953 let track = TrackInfo {
954 id,
955 path: path.clone(),
956 codec: info.codec.clone(),
957 sample_rate: info.sample_rate,
958 bit_depth: info.bit_depth,
959 bitrate_kbps: info.bitrate_kbps,
960 channels: info.channels,
961 duration_ms: info.duration_ms,
962 };
963 self.shared_state.set_position_ms(seek_ms);
966 self.on_track_changed(id, seek_ms);
967 log::info!(
968 "{}: {} ({:?}) — {} {}Hz/{}ch, {}ms{}",
969 if streaming { "streaming" } else { "playing" },
970 path.display(),
971 id,
972 info.codec,
973 info.sample_rate,
974 info.channels,
975 info.duration_ms,
976 if seek_ms > 0 {
977 format!(" @{}ms", seek_ms)
978 } else {
979 String::new()
980 }
981 );
982
983 let (producer, consumer) = rtrb::RingBuffer::new(RING_BUFFER_SIZE);
984 self.timeline.reset();
985 let lookahead = Arc::new(parking_lot::Mutex::new(Vec::new()));
986 let next_track = self.decode_cursor(id, lookahead.clone());
987 let processing = self.processing();
990 self.session += 1;
991 let session = self.session;
992 let finish_tx = self.commands.tx.clone();
993 let decode_handle = buffer::start_decode(
994 first,
995 producer,
996 seek_ms,
997 move || {
998 let (next_id, next_path) = next_track()?;
999 Some(buffer::SourceEntry::from_file(next_id, next_path))
1000 },
1001 self.timeline.clone(),
1002 Some(self.viz_buffer.clone()),
1003 processing,
1004 move || {
1005 finish_tx.send(PlayerCommand::DecodeFinished(session)).ok();
1006 },
1007 )?;
1008
1009 let (engine, rate_watch) = self.create_engine_for(&info, consumer)?;
1010 let dsp = self
1012 .dsp
1013 .as_ref()
1014 .and_then(|c| c.setup.as_ref())
1015 .map(|s| s.status(info.sample_rate));
1016 if start == Run::Playing {
1019 engine.start()?;
1020 }
1021 self.transport = Transport::Loaded(Session {
1022 track,
1023 run: start,
1024 lookahead,
1025 output: Output::Local(Local {
1026 engine,
1027 decode_handle,
1028 stream,
1029 _rate_watch: rate_watch,
1030 dsp,
1031 }),
1032 });
1033 Ok(())
1034 }
1035
1036 fn decode_cursor(
1046 &self,
1047 id: QueueItemId,
1048 steps: Arc<parking_lot::Mutex<Vec<state::Lookahead>>>,
1049 ) -> impl Fn() -> Option<(QueueItemId, PathBuf)> + Send + 'static {
1050 let state = self.shared_state.clone();
1051 let timeline = self.timeline.clone();
1052 move || {
1053 let mut steps = steps.lock();
1054 let current = match steps.last() {
1055 None => id,
1056 Some(step) => step.chosen.as_ref()?.0,
1057 };
1058 let mut step = state.lookahead_after(current)?;
1059 step.boundary = timeline.boundary_count();
1060 let next = step.chosen.clone();
1061 steps.push(step);
1062 next
1063 }
1064 }
1065
1066 fn probe_stream_for_playback(
1081 &self,
1082 id: QueueItemId,
1083 path: &Path,
1084 bytes_written: Arc<crate::remote::downloads::ByteFeed>,
1085 total: u64,
1086 ) {
1087 let path = path.to_path_buf();
1088 let tx = self.commands.tx.clone();
1089 let hint = hint_for(&path);
1090
1091 let status = {
1096 let downloading = self.stream_status_fn(id);
1097 let state = self.shared_state.clone();
1098 Arc::new(move || {
1099 if state.is_cursor(id) {
1100 downloading()
1101 } else {
1102 streaming::StreamStatus::Failed
1103 }
1104 }) as Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync>
1105 };
1106
1107 let spawned = thread::Builder::new()
1108 .name("koan-stream-probe".into())
1109 .spawn(move || {
1110 let attempt = |mode, wait: bool| {
1117 let open = if wait {
1118 streaming::PartialFileSource::open(
1119 &path,
1120 bytes_written.clone(),
1121 total,
1122 status.clone(),
1123 mode,
1124 )
1125 } else {
1126 streaming::PartialFileSource::open_for_probe(
1127 &path,
1128 bytes_written.clone(),
1129 total,
1130 status.clone(),
1131 mode,
1132 )
1133 };
1134 open.map_err(buffer::DecodeError::Io).and_then(|source| {
1135 let mss = symphonia::core::io::MediaSourceStream::new(
1136 Box::new(source),
1137 Default::default(),
1138 );
1139 buffer::probe_source(mss, &hint)
1140 })
1141 };
1142
1143 let info = match attempt(streaming::ProbeMode::Full, false) {
1147 Ok(info) => Some((info, streaming::ProbeMode::Full)),
1148 Err(e) => {
1149 log::info!(
1155 "stream probe: {} needs more than has arrived ({}), opening without a length",
1156 path.display(),
1157 e
1158 );
1159 let lengthless = lengthless_mode_for(&path);
1160 attempt(lengthless, true)
1161 .ok()
1162 .map(|info| (info, lengthless))
1163 }
1164 };
1165
1166 match info {
1167 Some((info, mode)) => {
1168 tx.send(PlayerCommand::StreamProbed {
1169 id,
1170 info: Box::new(info),
1171 mode,
1172 })
1173 .ok();
1174 }
1175 None => log::info!(
1178 "stream probe: {} cannot start early, waiting for the download",
1179 path.display()
1180 ),
1181 }
1182 });
1183
1184 if let Err(e) = spawned {
1185 log::warn!("stream probe: could not spawn for {:?}: {}", id, e);
1186 }
1187 }
1188
1189 fn stream_probed(
1192 &mut self,
1193 id: QueueItemId,
1194 info: buffer::StreamInfo,
1195 mode: streaming::ProbeMode,
1196 ) {
1197 let Some(waiting) = self.waiting().filter(|w| self.may_stream(w, id)) else {
1200 return;
1201 };
1202 if !self.shared_state.is_cursor(id) {
1203 return;
1204 }
1205
1206 match self.shared_state.item_playback_source(id) {
1207 Some(PlaybackSource::Ready(path)) => {
1209 if let Err(e) = self.start_playback(id, &path, 0, waiting.start) {
1210 log::error!("stream probe: playback failed: {}", e);
1211 }
1212 }
1213 Some(PlaybackSource::Streaming {
1214 path,
1215 bytes_written,
1216 total,
1217 }) => {
1218 let source = StreamSource {
1219 path,
1220 bytes_written,
1221 total,
1222 mode,
1223 };
1224 if let Err(e) =
1225 self.open_session(id, Source::Stream(source), Some(info), 0, waiting.start)
1226 {
1227 log::error!("stream probe: streaming playback failed: {}", e);
1228 }
1229 }
1230 None => {}
1231 }
1232 }
1233
1234 fn stream_status_fn(
1238 &self,
1239 id: QueueItemId,
1240 ) -> Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync> {
1241 let state = self.shared_state.clone();
1245 Arc::new(move || match state.item_state(id) {
1246 Some(ItemState::Ready) => streaming::StreamStatus::Complete,
1247 Some(ItemState::Failed(_)) => streaming::StreamStatus::Failed,
1248 _ => streaming::StreamStatus::Downloading,
1249 })
1250 }
1251
1252 pub fn seek(&mut self, position_ms: u64) {
1259 let Some(id) = self.session().map(|s| s.track.id) else {
1260 return;
1261 };
1262 let seekable = self.shared_state.seekable_ms();
1264 if seekable == 0 {
1265 log::debug!("seek declined: {:?} is not seekable yet", id);
1269 return;
1270 }
1271 let ceiling = seekable.min(
1272 self.shared_state
1273 .duration_ms()
1274 .saturating_sub(SEEK_END_GUARD_MS),
1275 );
1276 let clamped = position_ms.min(ceiling);
1277
1278 if self.seek_on_renderer(clamped) {
1281 return;
1282 }
1283 if let Err(e) = self.restart_current(clamped) {
1284 log::error!("seek failed: {}", e);
1285 }
1286 }
1287
1288 fn restart_current(&mut self, position_ms: u64) -> Result<(), PlayerError> {
1296 let Some(session) = self.session() else {
1297 return Ok(());
1298 };
1299 let (info, start) = (session.track.clone(), session.run);
1300
1301 match self.shared_state.item_playback_source(info.id) {
1302 Some(PlaybackSource::Streaming { .. }) if self.renderer.is_some() => {
1305 self.park(info.id, position_ms, start);
1306 return Ok(());
1307 }
1308 Some(PlaybackSource::Streaming {
1309 path,
1310 bytes_written,
1311 total,
1312 }) => {
1313 let known = buffer::StreamInfo {
1317 codec: info.codec.clone(),
1318 sample_rate: info.sample_rate,
1319 channels: info.channels,
1320 bit_depth: info.bit_depth,
1321 bitrate_kbps: info.bitrate_kbps,
1322 duration_ms: info.duration_ms,
1323 };
1324 let source = StreamSource {
1325 path,
1326 bytes_written,
1327 total,
1328 mode: self.stream_mode,
1329 };
1330 self.open_session(
1331 info.id,
1332 Source::Stream(source),
1333 Some(known),
1334 position_ms,
1335 start,
1336 )?;
1337 }
1338 Some(PlaybackSource::Ready(path)) => {
1339 self.start_playback(info.id, &path, position_ms, start)?;
1340 }
1341 None => return Ok(()),
1342 }
1343
1344 self.report(match start {
1345 Run::Playing => PlaybackReportState::Playing,
1346 Run::Paused => PlaybackReportState::Paused,
1347 });
1348 Ok(())
1349 }
1350
1351 pub fn next_track(&mut self) {
1353 match self.shared_state.advance_cursor_loadable() {
1354 Some(id) => self.play(id),
1355 None => {
1356 log::info!("no more tracks in playlist");
1357 self.stop_playback_and_clear_state();
1358 }
1359 }
1360 }
1361
1362 pub fn prev_track(&mut self) {
1364 match self.shared_state.retreat_cursor() {
1365 Some((id, _)) => self.play(id),
1366 None => {
1367 if let Err(e) = self.restart_current(0) {
1369 log::error!("restart failed: {}", e);
1370 }
1371 }
1372 }
1373 }
1374
1375 pub fn pause(&mut self) {
1381 let fade = crate::config::Config::cached().playback.fade_on_pause;
1382 match &mut self.transport {
1383 Transport::Idle => return,
1384 Transport::Waiting(waiting) => {
1385 waiting.start = Run::Paused;
1386 return;
1387 }
1388 Transport::Loaded(session) => {
1389 if let Output::Local(local) = &session.output {
1390 if fade {
1391 local.engine.fade_out();
1392 } else if let Err(e) = local.engine.stop() {
1393 log::error!("pause failed: {}", e);
1394 return;
1395 }
1396 }
1397 session.run = Run::Paused;
1398 }
1399 }
1400 self.pause_renderer();
1401 self.report(PlaybackReportState::Paused);
1402 }
1403
1404 pub fn resume(&mut self) {
1411 self.answer_silence();
1412 let session = match &mut self.transport {
1413 Transport::Idle => {
1414 if let Some(id) = self.shared_state.cursor() {
1415 self.play(id);
1416 }
1417 return;
1418 }
1419 Transport::Waiting(waiting) => {
1420 let waiting = *waiting;
1421 self.cue(waiting.id, waiting.position_ms, Run::Playing);
1422 return;
1423 }
1424 Transport::Loaded(session) => session,
1425 };
1426 let Output::Local(local) = &session.output else {
1427 self.resume_renderer();
1428 return;
1429 };
1430 let engine = &local.engine;
1431 let resumed = if engine.is_running() || engine.is_silent() {
1432 self.lead_in_ends = None;
1433 engine.fade_in()
1434 } else {
1435 engine.start()
1436 };
1437 if let Err(e) = resumed {
1438 log::error!("resume failed: {}", e);
1439 return;
1440 }
1441 session.run = Run::Playing;
1442 self.wake_analyzer();
1443 self.report(PlaybackReportState::Playing);
1444 }
1445
1446 fn wake_analyzer(&self) {
1453 self.viz_snapshot.wake();
1454 }
1455
1456 pub fn stop(&mut self) {
1458 self.shared_state.clear_playlist();
1459 self.stop_playback_and_clear_state();
1460 }
1461
1462 fn stop_engine(&mut self) {
1468 let playback = match std::mem::replace(&mut self.transport, Transport::Idle) {
1469 Transport::Loaded(playback) => playback,
1470 other => {
1471 self.transport = other;
1472 return;
1473 }
1474 };
1475 self.session += 1;
1476 self.bank_listening();
1477 let local = match playback.output {
1478 Output::Local(local) => local,
1479 Output::Renderer(play) => {
1480 self.halt_renderer(*play);
1481 self.answer_silence();
1482 return;
1483 }
1484 };
1485 let Local {
1486 engine,
1487 mut decode_handle,
1488 stream,
1489 ..
1490 } = local;
1491
1492 let _ = engine.stop();
1493 decode_handle.signal_stop();
1496 if let Some(stream) = stream {
1497 stream.abandon();
1498 }
1499 decode_handle.stop();
1500 drop(engine);
1501 self.answer_silence();
1502 }
1503
1504 fn answer_silence(&mut self) {
1506 let position_ms = self.shared_state.position_ms();
1507 for reply in self.silence_waiters.drain(..) {
1508 let _ = reply.send(position_ms);
1509 }
1510 }
1511
1512 fn stop_playback_and_clear_state(&mut self) {
1514 self.forget_waiting();
1515 self.report(PlaybackReportState::Stopped);
1516 self.finish_play();
1517 self.stop_engine();
1518 self.timeline.reset();
1519 self.shared_state.set_position_ms(0);
1520 }
1521
1522 pub fn remove_from_playlist(&mut self, id: QueueItemId) {
1529 let was_cursor = self.shared_state.is_cursor(id);
1530 let resume_after = was_cursor
1531 .then(|| self.shared_state.item_before(id))
1532 .flatten();
1533 self.shared_state.remove_item(id);
1534 if was_cursor {
1535 self.shared_state.set_cursor(resume_after);
1536 let next = self.shared_state.advance_cursor_loadable();
1537 self.carry_on(next, self.intent());
1538 }
1539 }
1540
1541 pub fn track_ready(&mut self, id: QueueItemId) {
1549 if !self.shared_state.is_cursor(id) {
1550 return;
1551 }
1552
1553 if let Some(waiting) = self.waiting().filter(|w| w.id == id) {
1554 log::info!("track_ready: opening {:?}", id);
1555 self.cue(id, waiting.position_ms, waiting.start);
1556 return;
1557 }
1558
1559 if self.session().is_some_and(|s| s.track.id == id) {
1560 log::info!(
1561 "track_ready: download complete while streaming {:?}, refreshing metadata",
1562 id
1563 );
1564 self.refresh_track_metadata(id);
1565 }
1566 }
1567
1568 pub fn track_stream_ready(&mut self, id: QueueItemId) {
1571 let Some(waiting) = self.waiting().filter(|w| self.may_stream(w, id)) else {
1572 return;
1573 };
1574 if !self.shared_state.is_cursor(id) {
1575 return;
1576 }
1577
1578 match self.shared_state.item_playback_source(id) {
1579 Some(PlaybackSource::Streaming {
1580 path,
1581 bytes_written,
1582 total,
1583 }) => {
1584 log::info!("track_stream_ready: probing partial file for {:?}", id);
1585 self.probe_stream_for_playback(id, &path, bytes_written, total);
1586 }
1587 Some(PlaybackSource::Ready(path)) => {
1588 log::info!(
1590 "track_stream_ready: track already ready, starting normal playback for {:?}",
1591 id
1592 );
1593 if let Err(e) = self.start_playback(id, &path, 0, waiting.start) {
1594 log::error!("track_stream_ready playback failed: {}", e);
1595 }
1596 }
1597 None => {} }
1599 }
1600
1601 fn refresh_track_metadata(&mut self, id: QueueItemId) {
1605 use crate::index::metadata;
1606
1607 let path = match self.shared_state.item_path_if_ready(id) {
1608 Some(p) => p,
1609 None => return,
1610 };
1611
1612 match metadata::read_metadata(&path) {
1613 Ok(meta) => {
1614 self.shared_state.update_item_metadata(
1615 id,
1616 meta.title,
1617 meta.artist,
1618 meta.album_artist.unwrap_or_default(),
1619 meta.album,
1620 meta.duration_ms.map(|d| d as u64),
1621 );
1622
1623 if let Transport::Loaded(session) = &mut self.transport
1632 && session.track.id == id
1633 {
1634 let current = &mut session.track;
1635 let probed = buffer::probe_file(&path).ok();
1636 let duration_ms = probed
1637 .as_ref()
1638 .map(|s| s.duration_ms)
1639 .filter(|d| *d > current.duration_ms)
1640 .unwrap_or(current.duration_ms);
1641 if duration_ms != current.duration_ms {
1642 log::info!(
1643 "track_ready: duration corrected {}ms → {}ms",
1644 current.duration_ms,
1645 duration_ms
1646 );
1647 }
1648 current.duration_ms = duration_ms;
1649 current.path = path.clone();
1650 }
1651
1652 self.shared_state.signal_metadata_refresh();
1654 log::info!("track_ready: metadata refreshed for {:?}", id);
1655 }
1656 Err(e) => {
1657 log::warn!("track_ready: metadata refresh failed for {:?}: {}", id, e);
1658 }
1659 }
1660 }
1661
1662 fn on_track_changed(&mut self, id: QueueItemId, position_ms: u64) {
1669 if let Some(f) = self.in_flight.as_mut().filter(|f| f.item == id) {
1670 f.jump(position_ms);
1671 f.boundary = 0;
1673 return;
1674 }
1675 self.begin_play(id, position_ms, 0);
1676 }
1677
1678 fn begin_play(&mut self, id: QueueItemId, position_ms: u64, boundary: usize) {
1683 self.finish_play();
1684 let track_id = self.shared_state.item_db_id(id);
1685 let mut flight = InFlight::new(id, track_id, position_ms);
1686 flight.boundary = boundary;
1687 self.in_flight = Some(flight);
1688 if let (Some(track_id), Some(recorder)) = (track_id, self.history.as_ref()) {
1689 recorder.record(PlayEvent::Started {
1690 track_id,
1691 position_ms,
1692 });
1693 }
1694 }
1695
1696 fn bank_listening(&mut self) {
1700 let renderer_at = self
1702 .shared_state
1703 .renderer_clock()
1704 .map(|_| self.shared_state.position_ms());
1705 if let Some(f) = self.in_flight.as_mut()
1706 && let Some(at) = self.timeline.position_in(f.boundary).or(renderer_at)
1707 {
1708 f.advance(at);
1709 }
1710 }
1711
1712 fn report(&self, state: PlaybackReportState) {
1716 let Some(track_id) = self.in_flight.as_ref().and_then(InFlight::track_id) else {
1717 return;
1718 };
1719 if let Some(recorder) = self.history.as_ref() {
1720 recorder.record(PlayEvent::Playback(PlaybackReport {
1721 track_id,
1722 state,
1723 position_ms: self.shared_state.position_ms(),
1724 }));
1725 }
1726 }
1727
1728 fn finish_play(&mut self) -> Option<PlayEvent> {
1731 self.bank_listening();
1732 let flight = self.in_flight.take()?;
1733 let event = PlayEvent::Finished {
1734 track_id: flight.track_id()?,
1735 listened_ms: flight.listened_ms(),
1736 };
1737 if let Some(recorder) = self.history.as_ref() {
1738 recorder.record(event);
1739 }
1740 Some(event)
1741 }
1742
1743 pub fn update_playback_state(&mut self) {
1747 if self.session().is_some()
1748 && self
1749 .lead_in_ends
1750 .is_some_and(|end| std::time::Instant::now() >= end)
1751 {
1752 self.lead_in_ends = None;
1756 self.shared_state.changed();
1757 }
1758
1759 if let Some(session) = self.session()
1760 && session.run == Run::Paused
1761 && let Some(engine) = session.engine()
1762 && engine.is_running()
1763 && engine.is_silent()
1764 {
1765 if let Err(e) = engine.stop() {
1766 log::error!("stopping after fade failed: {}", e);
1767 }
1768 self.answer_silence();
1769 }
1770
1771 self.follow_playhead();
1772 self.renderer_tick();
1773 self.publish();
1774 }
1775
1776 fn follow_playhead(&mut self) {
1781 if self.session().is_none() {
1782 return;
1783 }
1784 let Some(playhead) = self.timeline.playhead() else {
1785 return;
1786 };
1787 let id = playhead.id;
1788 if self
1789 .in_flight
1790 .as_ref()
1791 .is_none_or(|f| f.item != id || f.boundary != playhead.boundary)
1792 {
1793 self.begin_play(id, playhead.position_ms, playhead.boundary);
1794 }
1795 let Transport::Loaded(session) = &mut self.transport else {
1796 return;
1797 };
1798 if session.track.id == id {
1799 return;
1800 }
1801 let Some((id, path, info, _)) = self.timeline.current_playback() else {
1802 return;
1803 };
1804 log::info!("timeline: now playing {:?}", id);
1805 session.track = TrackInfo {
1806 id,
1807 path,
1808 codec: info.codec,
1809 sample_rate: info.sample_rate,
1810 bit_depth: info.bit_depth,
1811 bitrate_kbps: info.bitrate_kbps,
1812 channels: info.channels,
1813 duration_ms: info.duration_ms,
1814 };
1815 self.shared_state.set_cursor(Some(id));
1816 }
1817
1818 pub fn track_failed(&mut self, id: QueueItemId) {
1826 let Some(waiting) = self.waiting().filter(|w| w.id == id) else {
1827 return;
1828 };
1829 if !self.shared_state.is_cursor(id) {
1830 return;
1831 }
1832 log::info!("track {:?} cannot load, moving on", id);
1833 let next = self.shared_state.advance_cursor_loadable();
1834 self.carry_on(next, Some(waiting.start));
1835 }
1836
1837 fn on_decode_finished(&mut self, session: u64) {
1843 if session != self.session || self.renderer_stream_finished() {
1844 return;
1845 }
1846 log::info!("decode finished, checking for next track");
1847 self.track_ended();
1848 }
1849
1850 fn track_ended(&mut self) {
1864 self.bank_listening();
1865 let ended = self.in_flight.as_ref().map(|f| f.item);
1866 let heard = self.in_flight.as_ref().is_some_and(|f| f.listened_ms() > 0);
1867 if self.mode.repeat != Repeat::Off && !heard {
1868 log::info!("track ended with nothing heard; not repeating it");
1869 self.stop_playback_and_clear_state();
1870 return;
1871 }
1872 let again = (self.mode.repeat == Repeat::One)
1873 .then(|| self.shared_state.cursor())
1874 .flatten()
1875 .filter(|id| {
1876 self.shared_state
1877 .item_state(*id)
1878 .is_some_and(|s| !matches!(s, ItemState::Failed(_)))
1879 });
1880 let next = again.or_else(|| self.shared_state.advance_cursor_loadable());
1881 if next.is_some() && next == ended {
1882 self.finish_play();
1883 }
1884 self.carry_on(next, self.intent());
1885 }
1886
1887 fn revoke_stale_lookahead(&mut self) {
1897 if self.session().is_none() {
1898 return;
1899 }
1900 self.update_playback_state();
1901 let Some(playhead) = self.timeline.playhead() else {
1902 return;
1903 };
1904 let position_ms = playhead.position_ms;
1905 let Some(session) = self.session().filter(|s| s.track.id == playhead.id) else {
1906 return;
1907 };
1908 let stale = {
1909 let steps = session.lookahead.lock();
1910 steps
1914 .iter()
1915 .filter(|step| step.boundary > playhead.boundary)
1916 .any(|step| !self.shared_state.still_follows(step))
1917 };
1918 if !stale {
1919 return;
1920 }
1921 log::info!("queue changed under the lookahead, restarting at {position_ms}ms");
1922 if let Err(e) = self.restart_current(position_ms) {
1923 log::error!("restart after a queue edit failed: {}", e);
1924 }
1925 }
1926
1927 fn snapshot_for_undo(
1931 &self,
1932 ids: &[QueueItemId],
1933 ) -> Vec<(Box<state::PlaylistItem>, Option<QueueItemId>)> {
1934 self.shared_state
1935 .items_before(ids)
1936 .into_iter()
1937 .filter_map(|(id, after)| Some((Box::new(self.shared_state.get_item(id)?), after)))
1938 .collect()
1939 }
1940
1941 fn push_undo(&mut self, entry: UndoEntry) {
1943 if let Some(ref mut batch) = self.batch_buffer {
1944 batch.push(entry);
1945 } else {
1946 self.undo_stack.push(entry);
1947 }
1948 }
1949
1950 pub fn process_command(&mut self, cmd: PlayerCommand) {
1952 if cmd.asks_to_play() && !matches!(cmd, PlayerCommand::Cue { .. })
1957 || matches!(
1958 cmd,
1959 PlayerCommand::UseRenderer(_)
1960 | PlayerCommand::SetOutputDevice(_)
1961 | PlayerCommand::ClearOutputDevice
1962 )
1963 {
1964 self.resume_renderer = false;
1965 }
1966 let edits_queue = matches!(
1967 cmd,
1968 PlayerCommand::AddToPlaylist(_)
1969 | PlayerCommand::InsertInPlaylist { .. }
1970 | PlayerCommand::RemoveFromPlaylist(_)
1971 | PlayerCommand::RemoveFromPlaylistBatch(_)
1972 | PlayerCommand::MoveInPlaylist { .. }
1973 | PlayerCommand::MoveItemsInPlaylist { .. }
1974 | PlayerCommand::ReorderPlaylist(_)
1975 | PlayerCommand::Undo
1976 | PlayerCommand::Redo
1977 | PlayerCommand::SetShuffle(_)
1978 | PlayerCommand::SetRepeat(_)
1979 | PlayerCommand::RestorePlayMode(_)
1980 );
1981 self.apply_command(cmd);
1982 if edits_queue {
1983 self.revoke_stale_lookahead();
1984 }
1985 self.publish();
1986 }
1987
1988 fn apply_command(&mut self, cmd: PlayerCommand) {
1989 match cmd {
1990 PlayerCommand::Play(id) => self.play(id),
1991 PlayerCommand::Cue {
1992 id,
1993 position_ms,
1994 play,
1995 } => self.cue(
1996 id,
1997 position_ms,
1998 if play { Run::Playing } else { Run::Paused },
1999 ),
2000 PlayerCommand::Pause => self.pause(),
2001 PlayerCommand::PauseAndReport(reply) => {
2002 self.pause();
2003 self.silence_waiters.push(reply);
2004 if self
2005 .session()
2006 .is_none_or(|p| p.engine().is_none_or(|e| !e.is_running()))
2007 {
2008 self.answer_silence();
2009 }
2010 }
2011 PlayerCommand::Resume => self.resume(),
2012 PlayerCommand::Stop => self.stop(),
2013 PlayerCommand::Seek(pos) => self.seek(pos),
2014 PlayerCommand::NextTrack => self.next_track(),
2015 PlayerCommand::PrevTrack => self.prev_track(),
2016 PlayerCommand::AddToPlaylist(items) => {
2017 let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
2018 let whole = self.shared_state.is_empty();
2019 self.shared_state.add_items(items);
2020 if whole
2022 && self.mode.shuffle
2023 && let Some(&first) = ids.first()
2024 {
2025 self.shared_state.shuffle_from(first);
2026 }
2027 self.push_undo(UndoEntry::Added { ids });
2028 }
2029 PlayerCommand::UpdatePaths(updates) => {
2030 self.shared_state.update_paths(&updates);
2031 if let Transport::Loaded(session) = &mut self.transport
2032 && let Some((_, new_path)) =
2033 updates.iter().find(|(id, _)| *id == session.track.id)
2034 {
2035 session.track.path = new_path.clone();
2036 }
2037 }
2038 PlayerCommand::InsertInPlaylist { items, after } => {
2039 let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
2040 self.shared_state.insert_items_after(items, after);
2041 self.push_undo(UndoEntry::Inserted { ids });
2042 }
2043 PlayerCommand::ClearPlaylist => self.clear_playlist(),
2044 PlayerCommand::ReplacePlaylist {
2045 items,
2046 start,
2047 position_ms,
2048 play,
2049 } => {
2050 if items.is_empty() {
2051 self.clear_playlist();
2052 return;
2053 }
2054 let start_id = items.get(start).unwrap_or(&items[0]).id;
2055 self.stop_playback_and_clear_state();
2058 let (old_items, cursor) = self.shared_state.replace_playlist(items);
2059 if self.mode.shuffle {
2060 self.shared_state.shuffle_from(start_id);
2061 }
2062 self.push_undo(UndoEntry::Replaced {
2063 items: old_items,
2064 cursor,
2065 });
2066 self.cue(
2067 start_id,
2068 position_ms,
2069 if play { Run::Playing } else { Run::Paused },
2070 );
2071 }
2072 PlayerCommand::RemoveFromPlaylist(id) => {
2073 let item = self.shared_state.get_item(id);
2074 let after = self.shared_state.item_before(id);
2075 self.remove_from_playlist(id);
2076 if let Some(item) = item {
2077 self.push_undo(UndoEntry::Removed {
2078 items: vec![(Box::new(item), after)],
2079 });
2080 }
2081 }
2082 PlayerCommand::RemoveFromPlaylistBatch(ids) => {
2083 let items_with_pos = self.snapshot_for_undo(&ids);
2087 let intent = self.intent();
2088 let resume_after = match self.shared_state.cursor() {
2089 Some(cursor) if ids.contains(&cursor) => {
2090 Some(self.shared_state.surviving_item_before(cursor, &ids))
2091 }
2092 _ => None,
2093 };
2094
2095 self.shared_state.remove_items(&ids);
2096
2097 if let Some(resume_after) = resume_after {
2098 self.shared_state.set_cursor(resume_after);
2099 let next = self.shared_state.advance_cursor_loadable();
2100 self.carry_on(next, intent);
2101 }
2102
2103 if !items_with_pos.is_empty() {
2104 self.push_undo(UndoEntry::Removed {
2105 items: items_with_pos,
2106 });
2107 }
2108 }
2109 PlayerCommand::MoveInPlaylist { id, target, after } => {
2110 let was_after = self.shared_state.item_before(id);
2111 self.shared_state.move_item(id, target, after);
2112 self.push_undo(UndoEntry::Moved { id, was_after });
2113 }
2114 PlayerCommand::MoveItemsInPlaylist { ids, target, after } => {
2115 let entries = self.shared_state.items_before(&ids);
2116 self.shared_state.move_items(&ids, target, after);
2117 self.push_undo(UndoEntry::MovedBatch { entries });
2118 }
2119 PlayerCommand::ReorderPlaylist(order) => {
2120 let entries = self.shared_state.items_before(&order);
2124 self.shared_state.reorder_to(&order);
2125 self.push_undo(UndoEntry::MovedBatch { entries });
2126 }
2127 PlayerCommand::TrackReady(id) => self.track_ready(id),
2128 PlayerCommand::DecodeFinished(session) => self.on_decode_finished(session),
2129 PlayerCommand::TrackQueued => {}
2132 PlayerCommand::TrackStreamReady(id) => self.track_stream_ready(id),
2133 PlayerCommand::StreamProbed { id, info, mode } => self.stream_probed(id, *info, mode),
2134 PlayerCommand::TrackFailed(id) => self.track_failed(id),
2135 PlayerCommand::CacheTracks(ids) => {
2136 if let Some(downloads) = &self.downloads {
2137 downloads.cache(ids);
2138 }
2139 }
2140 PlayerCommand::Undo => self.execute_undo(),
2141 PlayerCommand::Redo => self.execute_redo(),
2142 PlayerCommand::BeginUndoBatch => {
2143 self.batch_buffer = Some(Vec::new());
2144 }
2145 PlayerCommand::EndUndoBatch => {
2146 if let Some(entries) = self.batch_buffer.take() {
2147 if entries.len() == 1 {
2148 self.undo_stack.push(entries.into_iter().next().unwrap());
2150 } else if !entries.is_empty() {
2151 self.undo_stack.push(UndoEntry::Batch(entries));
2152 }
2153 }
2154 }
2155 PlayerCommand::SetOutputDevice(name) => self.set_output_device(name),
2156 PlayerCommand::RestartOutput => {
2157 log::info!("restarting audio output");
2158 self.restart_on_current_track();
2159 }
2160 PlayerCommand::ClearOutputDevice => self.clear_output_device(),
2161 PlayerCommand::ReloadDsp => self.reload_dsp(),
2162 PlayerCommand::UseRenderer(connection) => {
2163 self.remember_renderer(connection.as_deref());
2164 self.use_renderer(connection);
2165 }
2166 PlayerCommand::ResumeRenderer(connection) => {
2167 if std::mem::take(&mut self.resume_renderer) && self.renderer.is_none() {
2168 log::info!(
2169 "upnp: back to {}, as last time",
2170 connection.session.renderer().name
2171 );
2172 self.use_renderer(Some(connection));
2173 } else {
2174 log::info!(
2175 "upnp: not going back to {}: playback or the output moved first",
2176 connection.session.renderer().name
2177 );
2178 }
2179 }
2180 PlayerCommand::SetRendererVolume(volume) => self.set_renderer_volume(volume),
2181 PlayerCommand::Renderer { session, event } => self.on_renderer_event(session, event),
2182 PlayerCommand::SetShuffle(on) => self.set_shuffle(on),
2183 PlayerCommand::SetRepeat(repeat) => self.mode.repeat = repeat,
2184 PlayerCommand::RestorePlayMode(mode) => self.mode = mode,
2185 }
2186 }
2187
2188 fn set_shuffle(&mut self, on: bool) {
2196 if self.mode.shuffle == on {
2197 return;
2198 }
2199 let order = self.shared_state.shuffle_order();
2200 if on {
2201 self.shared_state.shuffle_after_cursor();
2202 } else {
2203 self.shared_state.unshuffle();
2204 }
2205 self.push_undo(UndoEntry::Shuffled {
2206 shuffle: self.mode.shuffle,
2207 order,
2208 });
2209 self.mode.shuffle = on;
2210 }
2211
2212 fn clear_playlist(&mut self) {
2216 self.stop_playback_and_clear_state();
2217 let (items, cursor) = self.shared_state.snapshot_playlist();
2218 self.shared_state.clear_playlist();
2219 self.push_undo(UndoEntry::Replaced { items, cursor });
2220 }
2221
2222 fn apply_entry(&mut self, entry: UndoEntry) -> UndoEntry {
2224 match entry {
2225 UndoEntry::Added { ids } | UndoEntry::Inserted { ids } => {
2226 let items_with_pos = self.snapshot_for_undo(&ids);
2228 self.shared_state.remove_items(&ids);
2229 UndoEntry::Removed {
2230 items: items_with_pos,
2231 }
2232 }
2233 UndoEntry::Removed { items } => {
2234 let mut ids = Vec::with_capacity(items.len());
2236 for (item, after) in items {
2237 ids.push(item.id);
2238 self.shared_state.insert_item_at(*item, after);
2239 }
2240 UndoEntry::Added { ids }
2241 }
2242 UndoEntry::Moved { id, was_after } => {
2243 let current_after = self.shared_state.item_before(id);
2244 self.shared_state.move_item_to(id, was_after);
2245 UndoEntry::Moved {
2246 id,
2247 was_after: current_after,
2248 }
2249 }
2250 UndoEntry::MovedBatch { entries } => {
2251 let ids: Vec<QueueItemId> = entries.iter().map(|(id, _)| *id).collect();
2252 let current_positions = self.shared_state.items_before(&ids);
2253 self.shared_state.move_items_to(&entries);
2254 UndoEntry::MovedBatch {
2255 entries: current_positions,
2256 }
2257 }
2258 UndoEntry::Replaced { items, cursor } => {
2259 let (current_items, current_cursor) = self.shared_state.snapshot_playlist();
2260 self.shared_state.restore_playlist(items, cursor);
2261 UndoEntry::Replaced {
2262 items: current_items,
2263 cursor: current_cursor,
2264 }
2265 }
2266 UndoEntry::Shuffled { shuffle, order } => {
2267 let inverse = UndoEntry::Shuffled {
2268 shuffle: self.mode.shuffle,
2269 order: self.shared_state.shuffle_order(),
2270 };
2271 self.shared_state.restore_shuffle_order(&order);
2272 self.mode.shuffle = shuffle;
2273 inverse
2274 }
2275 UndoEntry::Batch(entries) => {
2276 let mut inverses: Vec<_> = entries
2278 .into_iter()
2279 .rev()
2280 .map(|e| self.apply_entry(e))
2281 .collect();
2282 inverses.reverse();
2283 UndoEntry::Batch(inverses)
2284 }
2285 }
2286 }
2287
2288 fn reconcile_playback(&mut self) {
2307 let orphaned = match &self.transport {
2308 Transport::Idle => false,
2309 Transport::Waiting(waiting) => self.shared_state.get_item(waiting.id).is_none(),
2310 Transport::Loaded(session) => self.shared_state.get_item(session.track.id).is_none(),
2311 };
2312 if !orphaned {
2313 return;
2314 }
2315 let cursor = self.shared_state.cursor();
2319 self.carry_on(cursor, self.intent());
2320 }
2321
2322 fn execute_undo(&mut self) {
2324 let Some(entry) = self.undo_stack.pop_undo() else {
2325 return;
2326 };
2327 let inverse = self.apply_entry(entry);
2328 self.undo_stack.push_redo(inverse);
2329 self.reconcile_playback();
2330 }
2331
2332 fn execute_redo(&mut self) {
2334 let Some(entry) = self.undo_stack.pop_redo() else {
2335 return;
2336 };
2337 let inverse = self.apply_entry(entry);
2338 self.undo_stack.push_undo_keep_redo(inverse);
2339 self.reconcile_playback();
2340 }
2341
2342 pub fn run(&mut self) {
2348 use crossbeam_channel::RecvTimeoutError;
2349
2350 let rx = self.commands.rx.clone();
2351 loop {
2352 let received = match self.next_wake() {
2353 Some(at) => rx.recv_deadline(at),
2354 None => rx.recv().map_err(|_| RecvTimeoutError::Disconnected),
2355 };
2356 match received {
2357 Ok(cmd) => self.process_command(cmd),
2358 Err(RecvTimeoutError::Timeout) => {}
2359 Err(RecvTimeoutError::Disconnected) => break,
2360 }
2361 self.update_playback_state();
2362 }
2363 self.stop();
2364 }
2365
2366 fn next_wake(&self) -> Option<std::time::Instant> {
2370 let session = self.session()?;
2371 let now = std::time::Instant::now();
2372 let Output::Local(local) = &session.output else {
2373 let next_track = (session.run == Run::Playing && self.streaming_to_renderer())
2375 .then(|| self.timeline.until_next_track())
2376 .flatten()
2377 .map(|left| now + left + BOUNDARY_SLACK);
2378 return match (next_track, self.renderer_deadline()) {
2379 (Some(a), Some(b)) => Some(a.min(b)),
2380 (a, b) => a.or(b),
2381 };
2382 };
2383 match session.run {
2384 Run::Playing => {
2385 let next_track = self
2386 .timeline
2387 .until_next_track()
2388 .map(|left| now + left + BOUNDARY_SLACK);
2389 match (next_track, self.lead_in_ends) {
2390 (Some(a), Some(b)) => Some(a.min(b)),
2391 (a, b) => a.or(b),
2392 }
2393 }
2394 Run::Paused if local.engine.is_running() => Some(now + FADE_CHECK),
2395 Run::Paused => None,
2396 }
2397 }
2398
2399 fn remember_renderer(&self, connection: Option<&crate::upnp::Connection>) {
2405 let renderer = connection.map(|c| c.session.renderer());
2406 if let Err(e) = crate::config::Config::persist(|cfg| {
2407 cfg.playback.renderer = renderer.map(|r| r.udn.clone());
2408 cfg.playback.renderer_name = renderer.map(|r| r.name.clone());
2409 }) {
2410 log::error!("failed to save the output: {e}");
2411 }
2412 }
2413
2414 pub fn spawn() -> (
2417 Arc<SharedPlayerState>,
2418 Arc<PlaybackTimeline>,
2419 Arc<VizSnapshot>,
2420 crossbeam_channel::Sender<PlayerCommand>,
2421 ) {
2422 Self::spawn_with(false)
2423 }
2424
2425 pub fn spawn_for_listening() -> (
2431 Arc<SharedPlayerState>,
2432 Arc<PlaybackTimeline>,
2433 Arc<VizSnapshot>,
2434 crossbeam_channel::Sender<PlayerCommand>,
2435 ) {
2436 Self::spawn_with(true)
2437 }
2438
2439 fn renderer_to_resume(listening: bool) -> Option<String> {
2442 listening
2443 .then(|| crate::config::Config::cached().playback.renderer.clone())
2444 .flatten()
2445 }
2446
2447 fn spawn_with(
2448 listening: bool,
2449 ) -> (
2450 Arc<SharedPlayerState>,
2451 Arc<PlaybackTimeline>,
2452 Arc<VizSnapshot>,
2453 crossbeam_channel::Sender<PlayerCommand>,
2454 ) {
2455 let mut player = Self::new();
2456 player.history = PlayRecorder::spawn();
2457 let state = player.shared_state();
2458 let timeline = player.timeline();
2459 let viz_snapshot = player.viz_snapshot();
2460 let tx = player.command_sender();
2461 player.downloads = Some(crate::remote::queue::DownloadQueue::spawn(
2464 tx.clone(),
2465 state.clone(),
2466 ));
2467
2468 if let Some(udn) = Self::renderer_to_resume(listening) {
2469 player.resume_renderer = true;
2470 crate::upnp::resume(udn, &tx);
2471 }
2472
2473 thread::Builder::new()
2474 .name("koan-player".into())
2475 .spawn(move || player.run())
2476 .expect("failed to spawn player thread");
2477
2478 (state, timeline, viz_snapshot, tx)
2479 }
2480}
2481
2482#[cfg(test)]
2483mod tests {
2484 #[test]
2487 fn a_headless_player_does_not_go_back_to_a_renderer() {
2488 crate::config::isolate_config_for_tests();
2489 crate::config::Config::persist(|cfg| {
2490 cfg.playback.renderer = Some("uuid:headless-test".into());
2491 })
2492 .unwrap();
2493 assert_eq!(Player::renderer_to_resume(false), None);
2494 }
2495
2496 #[test]
2497 fn a_download_in_progress_is_known_by_its_own_extension() {
2498 use std::path::Path;
2499 assert_eq!(
2501 lengthless_mode_for(Path::new("/c/t.m4a.part")),
2502 streaming::ProbeMode::LengthlessWholeEnd
2503 );
2504 assert_eq!(
2505 lengthless_mode_for(Path::new("/c/t.OPUS.part")),
2506 streaming::ProbeMode::LengthlessWholeEnd
2507 );
2508 assert_eq!(
2509 lengthless_mode_for(Path::new("/c/t.flac.part")),
2510 streaming::ProbeMode::Lengthless
2511 );
2512 assert_eq!(
2513 media_extension(Path::new("/c/t.m4a.part")).as_deref(),
2514 Some("m4a")
2515 );
2516 assert_eq!(
2517 media_extension(Path::new("/c/t.mp3")).as_deref(),
2518 Some("mp3")
2519 );
2520 }
2521
2522 use super::*;
2523 use state::PlaylistItem;
2524 use std::path::PathBuf;
2525 use std::sync::atomic::{AtomicU64, Ordering};
2526
2527 fn make_item(title: &str) -> PlaylistItem {
2528 PlaylistItem {
2529 playlist_entry_id: None,
2530 id: QueueItemId::new(),
2531 db_id: None,
2532 path: PathBuf::from(format!("/music/{title}.flac")),
2533 title: title.to_string(),
2534 artist: String::new(),
2535 album_artist: String::new(),
2536 album: String::new(),
2537 year: None,
2538 codec: None,
2539 track_number: None,
2540 disc: None,
2541 duration_ms: None,
2542 state: ItemState::Ready,
2543 pre_shuffle: None,
2544 }
2545 }
2546
2547 pub(super) fn playlist_ids(player: &Player) -> Vec<QueueItemId> {
2548 let (items, _) = player.shared_state.snapshot_playlist();
2549 items.iter().map(|i| i.id).collect()
2550 }
2551
2552 fn playlist_titles(player: &Player) -> Vec<String> {
2553 let (items, _) = player.shared_state.snapshot_playlist();
2554 items.iter().map(|i| i.title.clone()).collect()
2555 }
2556
2557 fn pending_item(title: &str) -> PlaylistItem {
2558 PlaylistItem {
2559 playlist_entry_id: None,
2560 state: ItemState::Pending,
2561 ..make_item(title)
2562 }
2563 }
2564
2565 fn test_session(id: QueueItemId, engine: Box<dyn AudioEngineHandle>) -> Session {
2567 Session {
2568 track: TrackInfo {
2569 id,
2570 path: PathBuf::from("/music/t.flac"),
2571 codec: String::new(),
2572 sample_rate: 44_100,
2573 bit_depth: None,
2574 bitrate_kbps: None,
2575 channels: 2,
2576 duration_ms: 1_000,
2577 },
2578 run: Run::Playing,
2579 lookahead: Default::default(),
2580 output: Output::Local(Local {
2581 engine,
2582 decode_handle: buffer::DecodeHandle::new_for_test(Default::default()),
2583 stream: None,
2584 _rate_watch: None,
2585 dsp: None,
2586 }),
2587 }
2588 }
2589
2590 fn pretend_playing(player: &mut Player, id: QueueItemId) {
2594 assert!(
2595 player.shared_state.get_item(id).is_some(),
2596 "item is in the queue"
2597 );
2598 player.stop_engine();
2599 player.transport = Transport::Loaded(test_session(
2600 id,
2601 Box::new(NullEngine {
2602 starts: Default::default(),
2603 running: Default::default(),
2604 lead_in: Default::default(),
2605 }),
2606 ));
2607 player.shared_state.set_cursor(Some(id));
2608 player.publish();
2609 }
2610
2611 fn playing_id(player: &Player) -> Option<QueueItemId> {
2612 player.shared_state.track_info().map(|t| t.id)
2613 }
2614
2615 fn seed(player: &mut Player, n: usize) -> Vec<QueueItemId> {
2617 let items: Vec<_> = (0..n).map(|i| make_item(&format!("t{i}"))).collect();
2618 let ids = items.iter().map(|i| i.id).collect();
2619 player.process_command(PlayerCommand::AddToPlaylist(items));
2620 ids
2621 }
2622
2623 fn listen(player: &mut Player, from_ms: u64, to_ms: u64) {
2627 let mut at = from_ms;
2628 if let Some(f) = player.in_flight.as_mut() {
2629 f.advance(at); }
2631 while at < to_ms {
2632 at = (at + 50).min(to_ms);
2633 if let Some(f) = player.in_flight.as_mut() {
2634 f.advance(at);
2635 }
2636 }
2637 }
2638
2639 fn start(player: &mut Player, track_id: i64) -> QueueItemId {
2640 let id = QueueItemId::new();
2641 player.on_track_changed(id, 0);
2642 player
2644 .in_flight
2645 .as_mut()
2646 .unwrap()
2647 .track_id_for_test(track_id);
2648 id
2649 }
2650
2651 #[test]
2652 fn a_gapless_transition_closes_the_outgoing_track_and_opens_the_next() {
2653 let mut player = Player::new();
2654 start(&mut player, 11);
2655 listen(&mut player, 0, 200_000);
2656
2657 let b = QueueItemId::new();
2658 player.on_track_changed(b, 0);
2659 let f = player
2660 .in_flight
2661 .as_ref()
2662 .expect("the next track is counting");
2663 assert_eq!(f.item, b);
2664 assert_eq!(f.listened_ms(), 0, "and starts from nothing");
2665 }
2666
2667 #[test]
2668 fn a_track_skipped_seconds_in_is_still_history() {
2669 let mut player = Player::new();
2670 start(&mut player, 7);
2671 listen(&mut player, 0, 2_000);
2672
2673 let event = player
2674 .finish_play()
2675 .expect("putting something on is a thing you did, however briefly");
2676 assert!(matches!(
2677 event,
2678 history::PlayEvent::Finished {
2679 track_id: 7,
2680 listened_ms: 2_000
2681 }
2682 ));
2683 }
2684
2685 #[test]
2686 fn a_track_is_closed_out_once() {
2687 let mut player = Player::new();
2688 start(&mut player, 7);
2689 listen(&mut player, 0, 200_000);
2690
2691 assert!(player.finish_play().is_some());
2692 assert!(player.finish_play().is_none());
2693 }
2694
2695 #[test]
2696 fn seeking_around_a_track_does_not_enter_it_twice() {
2697 let mut player = Player::new();
2698 let id = start(&mut player, 7);
2699 listen(&mut player, 0, 120_000);
2700
2701 player.on_track_changed(id, 30_000);
2703 assert_eq!(
2704 player.in_flight.as_ref().unwrap().listened_ms(),
2705 120_000,
2706 "the seek kept the count rather than restarting it"
2707 );
2708 listen(&mut player, 30_000, 40_000);
2709
2710 let Some(history::PlayEvent::Finished { listened_ms, .. }) = player.finish_play() else {
2711 panic!("still one play");
2712 };
2713 assert_eq!(listened_ms, 130_000);
2714 assert!(player.finish_play().is_none());
2715 }
2716
2717 #[test]
2718 fn a_track_that_is_not_in_the_library_is_not_recorded() {
2719 let mut player = Player::new();
2720 let id = QueueItemId::new();
2721 player.on_track_changed(id, 0);
2722 listen(&mut player, 0, 200_000);
2723 assert!(player.finish_play().is_none());
2724 }
2725
2726 #[test]
2727 fn stopping_closes_out_what_was_heard() {
2728 let mut player = Player::new();
2729 start(&mut player, 7);
2730 listen(&mut player, 0, 150_000);
2731
2732 player.stop_playback_and_clear_state();
2733 assert!(player.in_flight.is_none(), "the stop consumed it");
2734 }
2735
2736 #[test]
2737 fn resume_with_nothing_loaded_plays_the_cursor() {
2738 let dir = tempfile::tempdir().unwrap();
2739 let path = dir.path().join("t.wav");
2740 crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
2741
2742 let mut player = Player::new();
2743 player.backend = Box::new(StuckBackend {
2744 rate: 8_000.0,
2745 asked: Default::default(),
2746 starts: Default::default(),
2747 });
2748 let item = PlaylistItem {
2749 db_id: Some(5),
2750 path,
2751 ..make_item("t")
2752 };
2753 let id = item.id;
2754 player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
2755 player.shared_state.set_cursor(Some(id));
2756 assert!(player.session().is_none());
2757
2758 player.process_command(PlayerCommand::Resume);
2759 assert!(player.session().is_some(), "the cursor's track started");
2760 assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
2761 player.process_command(PlayerCommand::Stop);
2762 }
2763
2764 #[test]
2765 fn a_player_with_nothing_coming_has_nothing_to_wake_for() {
2766 let dir = tempfile::tempdir().unwrap();
2767 let path = dir.path().join("t.wav");
2768 crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
2769
2770 let mut player = Player::new();
2771 player.backend = Box::new(StuckBackend {
2772 rate: 8_000.0,
2773 asked: Default::default(),
2774 starts: Default::default(),
2775 });
2776 let item = PlaylistItem {
2777 path,
2778 ..make_item("t")
2779 };
2780 let id = item.id;
2781 player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
2782 assert_eq!(player.next_wake(), None, "stopped");
2783
2784 player.process_command(PlayerCommand::Play(id));
2785 assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
2786 assert_eq!(
2787 player.next_wake(),
2788 None,
2789 "playing, with no track queued after it"
2790 );
2791
2792 if let Transport::Loaded(session) = &mut player.transport {
2793 session.engine().unwrap().stop().unwrap();
2794 session.run = Run::Paused;
2795 }
2796 assert_eq!(player.next_wake(), None, "paused");
2797 player.process_command(PlayerCommand::Stop);
2798 }
2799
2800 #[test]
2801 fn a_track_cued_paused_never_starts_the_output() {
2802 let dir = tempfile::tempdir().unwrap();
2803 let path = dir.path().join("t.wav");
2804 crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
2805
2806 let starts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
2807 let mut player = Player::new();
2808 player.backend = Box::new(StuckBackend {
2809 rate: 8_000.0,
2810 asked: Default::default(),
2811 starts: starts.clone(),
2812 });
2813 let item = PlaylistItem {
2814 path,
2815 ..make_item("t")
2816 };
2817 let id = item.id;
2818 player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
2819
2820 player.process_command(PlayerCommand::Cue {
2821 id,
2822 position_ms: 2_000,
2823 play: false,
2824 });
2825 assert!(player.session().is_some(), "loaded");
2826 assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
2827 loop {
2831 match player
2832 .commands
2833 .rx
2834 .recv_timeout(std::time::Duration::from_secs(5))
2835 {
2836 Ok(PlayerCommand::TrackQueued) => break,
2837 Ok(_) => {}
2838 Err(e) => panic!("the decoder never queued the track: {e}"),
2839 }
2840 }
2841 let at = player.shared_state.position_ms();
2842 assert!((1_750..=2_000).contains(&at), "cued at {at}ms");
2843 assert_eq!(starts.load(Ordering::Relaxed), 0, "not a sample let out");
2844
2845 player.process_command(PlayerCommand::Seek(4_000));
2847 assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
2848 assert_eq!(starts.load(Ordering::Relaxed), 0);
2849
2850 player.process_command(PlayerCommand::Resume);
2851 assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
2852 assert_eq!(starts.load(Ordering::Relaxed), 1);
2853 player.process_command(PlayerCommand::Stop);
2854 }
2855
2856 fn downloading_wav(dir: &Path) -> (Player, QueueItemId, Arc<std::sync::atomic::AtomicUsize>) {
2859 let path = dir.join("t.wav");
2860 crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
2861 let starts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
2862 let mut player = Player::new();
2863 player.backend = Box::new(StuckBackend {
2864 rate: 8_000.0,
2865 asked: Default::default(),
2866 starts: starts.clone(),
2867 });
2868 let item = PlaylistItem {
2869 path,
2870 state: ItemState::Pending,
2871 ..make_item("t")
2872 };
2873 let id = item.id;
2874 player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
2875 (player, id, starts)
2876 }
2877
2878 fn land(player: &mut Player, id: QueueItemId) {
2881 player.shared_state.update_item_state(id, ItemState::Ready);
2882 player.process_command(PlayerCommand::TrackReady(id));
2883 }
2884
2885 #[test]
2886 fn track_ready_never_resurrects_an_item_that_failed() {
2887 let mut player = Player::new();
2888 let item = PlaylistItem {
2889 state: ItemState::Failed("gone".into()),
2890 ..make_item("t")
2891 };
2892 let id = item.id;
2893 player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
2894
2895 player.process_command(PlayerCommand::TrackReady(id));
2896
2897 assert!(matches!(
2898 player.shared_state.get_item(id).map(|i| i.state),
2899 Some(ItemState::Failed(_))
2900 ));
2901 }
2902
2903 fn await_queued(player: &Player) {
2904 loop {
2905 match player
2906 .commands
2907 .rx
2908 .recv_timeout(std::time::Duration::from_secs(5))
2909 {
2910 Ok(PlayerCommand::TrackQueued) => return,
2911 Ok(_) => {}
2912 Err(e) => panic!("the decoder never queued the track: {e}"),
2913 }
2914 }
2915 }
2916
2917 #[test]
2918 fn a_cue_on_a_downloading_track_opens_at_its_position_once_it_lands() {
2919 let dir = tempfile::tempdir().unwrap();
2920 let (mut player, id, starts) = downloading_wav(dir.path());
2921
2922 player.process_command(PlayerCommand::Cue {
2923 id,
2924 position_ms: 6_000,
2925 play: true,
2926 });
2927 assert!(player.session().is_none(), "nothing opened early");
2928 assert_eq!(player.shared_state.playback_state(), PlaybackState::Stopped);
2929 assert_eq!(starts.load(Ordering::Relaxed), 0);
2930
2931 land(&mut player, id);
2932 assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
2933 assert_eq!(player.playback_starts, 1, "opened once, at the position");
2934 assert_eq!(starts.load(Ordering::Relaxed), 1);
2935 await_queued(&player);
2936 let at = player.shared_state.position_ms();
2937 assert!((5_750..=6_000).contains(&at), "opened at {at}ms");
2938 player.process_command(PlayerCommand::Stop);
2939 }
2940
2941 #[test]
2942 fn a_paused_cue_on_a_downloading_track_lands_paused() {
2943 let dir = tempfile::tempdir().unwrap();
2944 let (mut player, id, starts) = downloading_wav(dir.path());
2945
2946 player.process_command(PlayerCommand::Cue {
2947 id,
2948 position_ms: 3_000,
2949 play: false,
2950 });
2951 land(&mut player, id);
2952 assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
2953 assert_eq!(starts.load(Ordering::Relaxed), 0, "not a sample let out");
2954 await_queued(&player);
2955 let at = player.shared_state.position_ms();
2956 assert!((2_750..=3_000).contains(&at), "cued at {at}ms");
2957 player.process_command(PlayerCommand::Stop);
2958 }
2959
2960 #[test]
2961 fn playing_something_else_forgets_a_waiting_cue() {
2962 let dir = tempfile::tempdir().unwrap();
2963 let (mut player, id, _) = downloading_wav(dir.path());
2964 let other = seed(&mut player, 1)[0];
2965
2966 player.process_command(PlayerCommand::Cue {
2967 id,
2968 position_ms: 6_000,
2969 play: true,
2970 });
2971 player.process_command(PlayerCommand::Play(other));
2972 player.process_command(PlayerCommand::Play(id));
2973 assert!(player.waiting().is_some_and(|w| w.position_ms == 0));
2974 player.process_command(PlayerCommand::Stop);
2975 }
2976
2977 #[test]
2978 fn a_paused_track_whose_download_lands_stays_paused_where_it_was() {
2979 let dir = tempfile::tempdir().unwrap();
2980 let (mut player, id, starts) = downloading_wav(dir.path());
2981 player.shared_state.update_item_state(id, ItemState::Ready);
2983 player.process_command(PlayerCommand::Cue {
2984 id,
2985 position_ms: 3_000,
2986 play: false,
2987 });
2988 await_queued(&player);
2989
2990 player.process_command(PlayerCommand::TrackReady(id));
2991 assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
2992 assert_eq!(player.playback_starts, 1, "not reopened");
2993 assert_eq!(starts.load(Ordering::Relaxed), 0, "not a sample let out");
2994 player.process_command(PlayerCommand::Stop);
2995 }
2996
2997 #[test]
2998 fn pausing_a_track_on_its_way_opens_it_paused() {
2999 let dir = tempfile::tempdir().unwrap();
3000 let (mut player, id, starts) = downloading_wav(dir.path());
3001
3002 player.process_command(PlayerCommand::Play(id));
3003 assert!(player.shared_state.is_waiting());
3004 assert!(player.shared_state.wants_to_play(), "a toggle pauses it");
3005 assert!(!player.shared_state.is_idle(), "adding tracks leaves it be");
3006 player.process_command(PlayerCommand::Pause);
3007 assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
3008 assert!(!player.shared_state.wants_to_play());
3009
3010 player.shared_state.update_item_state(id, ItemState::Ready);
3011 player.process_command(PlayerCommand::TrackReady(id));
3012 assert!(player.session().is_some(), "loaded");
3013 assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
3014 assert_eq!(starts.load(Ordering::Relaxed), 0, "not a sample let out");
3015 player.process_command(PlayerCommand::Stop);
3016 }
3017
3018 #[test]
3019 fn a_cue_paused_and_resumed_on_its_way_keeps_its_position() {
3020 let dir = tempfile::tempdir().unwrap();
3021 let (mut player, id, starts) = downloading_wav(dir.path());
3022
3023 player.process_command(PlayerCommand::Cue {
3024 id,
3025 position_ms: 6_000,
3026 play: true,
3027 });
3028 player.process_command(PlayerCommand::Pause);
3029 player.process_command(PlayerCommand::Resume);
3030 assert!(player.session().is_none(), "still on its way");
3031 assert_eq!(player.shared_state.playback_state(), PlaybackState::Stopped);
3032
3033 player.shared_state.update_item_state(id, ItemState::Ready);
3034 player.process_command(PlayerCommand::TrackReady(id));
3035 assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
3036 assert_eq!(starts.load(Ordering::Relaxed), 1);
3037 await_queued(&player);
3038 let at = player.shared_state.position_ms();
3039 assert!((5_750..=6_000).contains(&at), "opened at {at}ms");
3040 player.process_command(PlayerCommand::Stop);
3041 }
3042
3043 #[test]
3044 fn a_restored_cursor_does_not_start_when_its_download_lands() {
3045 let dir = tempfile::tempdir().unwrap();
3046 let (mut player, id, starts) = downloading_wav(dir.path());
3047 player.shared_state.set_cursor(Some(id));
3048
3049 player.shared_state.update_item_state(id, ItemState::Ready);
3050 player.process_command(PlayerCommand::TrackReady(id));
3051 assert!(player.session().is_none());
3052 assert_eq!(starts.load(Ordering::Relaxed), 0);
3053 }
3054
3055 #[test]
3056 fn replacing_the_queue_keeps_a_transfer_both_queues_want() {
3057 let pending = |title: &str| PlaylistItem {
3063 db_id: Some(7),
3064 state: ItemState::Pending,
3065 ..make_item(title)
3066 };
3067 let mut player = Player::new();
3068 let old = pending("old");
3069 let old_id = old.id;
3070 player.process_command(PlayerCommand::AddToPlaylist(vec![old]));
3071 let store = player.shared_state.downloads().clone();
3072 store.claim(7, Some(old_id));
3073
3074 let before = player.shared_state.pending_version();
3075 player.process_command(PlayerCommand::ReplacePlaylist {
3076 items: vec![pending("again")],
3077 start: 0,
3078 position_ms: 0,
3079 play: false,
3080 });
3081 assert_eq!(
3082 player.shared_state.pending_version(),
3083 before + 1,
3084 "one change, so no reader can see the playlist between two"
3085 );
3086
3087 store.resync(&player.shared_state.pending_downloads());
3089 assert!(!store.abandoned(7));
3090 }
3091
3092 #[test]
3093 fn a_hand_off_opens_paused_at_its_position_in_one_command() {
3094 let dir = tempfile::tempdir().unwrap();
3095 let path = dir.path().join("t.wav");
3096 crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
3097 let starts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
3098 let mut player = Player::new();
3099 player.backend = Box::new(StuckBackend {
3100 rate: 8_000.0,
3101 asked: Default::default(),
3102 starts: starts.clone(),
3103 });
3104 seed(&mut player, 2);
3105 let item = PlaylistItem {
3106 path,
3107 ..make_item("t")
3108 };
3109 let id = item.id;
3110
3111 player.process_command(PlayerCommand::ReplacePlaylist {
3112 items: vec![make_item("before"), item],
3113 start: 1,
3114 position_ms: 4_000,
3115 play: false,
3116 });
3117 assert_eq!(
3118 playlist_titles(&player),
3119 vec!["before", "t"],
3120 "replaced, not added to"
3121 );
3122 assert_eq!(player.shared_state.cursor(), Some(id));
3123 assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
3124 assert_eq!(player.playback_starts, 1, "opened once, at the position");
3125 assert_eq!(starts.load(Ordering::Relaxed), 0, "not a sample let out");
3126 await_queued(&player);
3127 let at = player.shared_state.position_ms();
3128 assert!((3_750..=4_000).contains(&at), "opened at {at}ms");
3129
3130 player.process_command(PlayerCommand::Undo);
3131 assert_eq!(playlist_titles(&player), vec!["t0", "t1"], "one undo step");
3132 player.process_command(PlayerCommand::Stop);
3133 }
3134
3135 #[test]
3136 fn a_decode_end_from_a_replaced_session_is_ignored() {
3137 let dir = tempfile::tempdir().unwrap();
3138 let mut player = Player::new();
3139 player.backend = Box::new(StuckBackend {
3140 rate: 8_000.0,
3141 asked: Default::default(),
3142 starts: Default::default(),
3143 });
3144 let items: Vec<_> = ["a", "b", "c"]
3145 .iter()
3146 .map(|name| {
3147 let path = dir.path().join(format!("{name}.wav"));
3148 crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
3149 PlaylistItem {
3150 path,
3151 ..make_item(name)
3152 }
3153 })
3154 .collect();
3155 let ids: Vec<_> = items.iter().map(|i| i.id).collect();
3156 player.process_command(PlayerCommand::AddToPlaylist(items));
3157
3158 player.process_command(PlayerCommand::Play(ids[0]));
3159 let stale = player.session;
3160 player.process_command(PlayerCommand::Play(ids[1]));
3161 player.process_command(PlayerCommand::DecodeFinished(stale));
3162
3163 assert_eq!(player.shared_state.cursor(), Some(ids[1]));
3164 assert_eq!(player.playback_starts, 2);
3165 player.process_command(PlayerCommand::Stop);
3166 }
3167
3168 fn queued_wavs(dir: &Path, names: &[&str]) -> (Player, Vec<QueueItemId>) {
3171 let (player, ids) = wavs_playing(dir, names, 1.0);
3172 for _ in &ids {
3173 await_queued(&player);
3174 }
3175 (player, ids)
3176 }
3177
3178 fn wavs_playing(dir: &Path, names: &[&str], seconds: f32) -> (Player, Vec<QueueItemId>) {
3181 wavs_in(dir, names, seconds, Repeat::Off)
3182 }
3183
3184 fn wavs_in(
3187 dir: &Path,
3188 names: &[&str],
3189 seconds: f32,
3190 repeat: Repeat,
3191 ) -> (Player, Vec<QueueItemId>) {
3192 let mut player = Player::new();
3193 player.backend = Box::new(StuckBackend {
3194 rate: 8_000.0,
3195 asked: Default::default(),
3196 starts: Default::default(),
3197 });
3198 let items: Vec<_> = names
3199 .iter()
3200 .zip(1..)
3201 .map(|(name, track)| {
3202 let path = dir.join(format!("{name}.wav"));
3203 crate::test_utils::generate_wav(&path, 8_000, 1, seconds, 16);
3204 PlaylistItem {
3205 path,
3206 db_id: Some(track),
3207 ..make_item(name)
3208 }
3209 })
3210 .collect();
3211 let ids: Vec<_> = items.iter().map(|i| i.id).collect();
3212 player.process_command(PlayerCommand::AddToPlaylist(items));
3213 player.process_command(PlayerCommand::SetRepeat(repeat));
3214 player.process_command(PlayerCommand::Play(ids[0]));
3215 (player, ids)
3216 }
3217
3218 fn queued_at_least(player: &Player, n: usize) -> Vec<QueueItemId> {
3220 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
3221 loop {
3222 let queued = player.timeline.queued_after_playhead();
3223 if queued.len() >= n {
3224 return queued;
3225 }
3226 assert!(
3227 std::time::Instant::now() < deadline,
3228 "the decoder queued {queued:?}, never {n}"
3229 );
3230 std::thread::yield_now();
3231 }
3232 }
3233
3234 #[test]
3237 fn repeating_the_queue_runs_the_last_track_into_the_first() {
3238 let dir = tempfile::tempdir().unwrap();
3239 let (mut player, ids) = wavs_in(dir.path(), &["a", "b"], 1.0, Repeat::Queue);
3240
3241 let queued = queued_at_least(&player, 3);
3242 assert_eq!(
3243 queued[..3],
3244 [ids[1], ids[0], ids[1]],
3245 "gapless, round again"
3246 );
3247 let steps = player.session().unwrap().lookahead.lock().clone();
3248 assert!(!steps[0].wrapped);
3249 let wrap = steps.iter().find(|s| s.after == ids[1]).unwrap();
3250 assert_eq!(wrap.next, Some(ids[0]));
3251 assert!(wrap.wrapped);
3252 assert_eq!(player.playback_starts, 1);
3253 player.process_command(PlayerCommand::Stop);
3254 }
3255
3256 #[test]
3257 fn turning_repeat_off_takes_back_a_wrap_already_queued() {
3258 let dir = tempfile::tempdir().unwrap();
3259 let (mut player, ids) = wavs_in(dir.path(), &["a"], 1.0, Repeat::Queue);
3260 assert_eq!(queued_at_least(&player, 1)[0], ids[0]);
3261
3262 player.process_command(PlayerCommand::SetRepeat(Repeat::Off));
3263 assert_eq!(player.playback_starts, 2, "restarted at the playhead");
3264 wait_for_lookahead(&player);
3265 assert!(player.timeline.queued_after_playhead().is_empty());
3266 player.process_command(PlayerCommand::Stop);
3267 }
3268
3269 #[test]
3270 fn repeating_one_track_queues_it_again_and_each_pass_is_a_play() {
3271 let dir = tempfile::tempdir().unwrap();
3272 let (mut player, ids) = wavs_in(dir.path(), &["a", "b"], 1.0, Repeat::One);
3273 let (recorder, events) = history::PlayRecorder::capture();
3274 player.history = Some(recorder);
3275 player.in_flight = None;
3278 player.on_track_changed(ids[0], 0);
3279
3280 assert_eq!(queued_at_least(&player, 2)[..2], [ids[0], ids[0]]);
3281 player
3283 .timeline
3284 .samples_played
3285 .store(8_400, Ordering::Relaxed);
3286 player.update_playback_state();
3287
3288 let flight = player.in_flight.as_ref().unwrap();
3289 assert_eq!((flight.item, flight.boundary), (ids[0], 1));
3290 assert_eq!(player.shared_state.cursor(), Some(ids[0]));
3291 let events: Vec<_> = events.try_iter().collect();
3292 let started = events
3293 .iter()
3294 .filter(|e| matches!(e, PlayEvent::Started { track_id: 1, .. }))
3295 .count();
3296 assert_eq!(started, 2, "two plays: {events:?}");
3297 assert!(
3298 events.contains(&PlayEvent::Finished {
3299 track_id: 1,
3300 listened_ms: 1_000
3301 }),
3302 "the first pass banked whole: {events:?}"
3303 );
3304 player.process_command(PlayerCommand::Stop);
3305 }
3306
3307 #[test]
3308 fn next_moves_on_from_a_track_repeating() {
3309 let dir = tempfile::tempdir().unwrap();
3310 let (mut player, ids) = wavs_in(dir.path(), &["a", "b"], 1.0, Repeat::One);
3311
3312 player.process_command(PlayerCommand::NextTrack);
3313 assert_eq!(playing_id(&player), Some(ids[1]));
3314 player.process_command(PlayerCommand::NextTrack);
3315 assert_eq!(
3316 playing_id(&player),
3317 Some(ids[0]),
3318 "round from the last, as repeating does"
3319 );
3320 player.process_command(PlayerCommand::Stop);
3321 }
3322
3323 #[test]
3324 fn a_repeat_of_nothing_heard_stops_rather_than_going_round() {
3325 let dir = tempfile::tempdir().unwrap();
3326 let (mut player, _) = wavs_in(dir.path(), &["a"], 1.0, Repeat::One);
3327
3328 player.process_command(PlayerCommand::DecodeFinished(player.session));
3331 assert!(matches!(player.transport, Transport::Idle));
3332 }
3333
3334 #[test]
3335 fn shuffle_on_and_off_again_puts_the_queue_back() {
3336 let mut player = Player::new();
3337 let ids = seed(&mut player, 20);
3338 pretend_playing(&mut player, ids[3]);
3339
3340 player.process_command(PlayerCommand::SetShuffle(true));
3341 let shuffled = playlist_ids(&player);
3342 assert!(player.shared_state.play_mode().shuffle);
3343 assert_eq!(
3344 shuffled[..4],
3345 ids[..4],
3346 "nothing up to the playing track moves"
3347 );
3348 assert_ne!(shuffled, ids);
3349 let mut sorted = shuffled.clone();
3350 sorted.sort_by_key(|id| ids.iter().position(|i| i == id));
3351 assert_eq!(sorted, ids, "the same items");
3352
3353 let extra = make_item("extra");
3354 let extra_id = extra.id;
3355 player.process_command(PlayerCommand::InsertInPlaylist {
3356 items: vec![extra],
3357 after: ids[3],
3358 });
3359 player.process_command(PlayerCommand::RemoveFromPlaylist(ids[10]));
3360 player.process_command(PlayerCommand::SetShuffle(false));
3361
3362 let mut expected = ids.clone();
3363 expected.remove(10);
3364 expected.insert(4, extra_id);
3365 assert_eq!(playlist_ids(&player), expected, "added since stays put");
3366 assert!(!player.shared_state.play_mode().shuffle);
3367 assert!(
3368 player
3369 .shared_state
3370 .shuffle_order()
3371 .iter()
3372 .all(|(_, pre)| pre.is_none())
3373 );
3374 }
3375
3376 #[test]
3377 fn a_queue_replaced_while_shuffled_plays_shuffled_from_its_start() {
3378 let mut player = Player::new();
3379 seed(&mut player, 3);
3380 player.process_command(PlayerCommand::SetShuffle(true));
3381
3382 let items: Vec<_> = (0..20).map(|i| make_item(&format!("n{i}"))).collect();
3383 let given: Vec<_> = items.iter().map(|i| i.id).collect();
3384 player.process_command(PlayerCommand::ReplacePlaylist {
3385 items,
3386 start: 5,
3387 position_ms: 0,
3388 play: false,
3389 });
3390 let shuffled = playlist_ids(&player);
3391 assert_eq!(shuffled[0], given[5], "the start first");
3392 assert_eq!(player.shared_state.cursor(), Some(given[5]));
3393 assert_ne!(shuffled, given);
3394
3395 player.process_command(PlayerCommand::SetShuffle(false));
3396 assert_eq!(
3397 playlist_ids(&player),
3398 given,
3399 "off gives the queue as it came"
3400 );
3401 }
3402
3403 #[test]
3404 fn a_queue_added_to_an_empty_one_while_shuffled_plays_shuffled() {
3405 let mut player = Player::new();
3406 player.process_command(PlayerCommand::SetShuffle(true));
3407 let given = seed(&mut player, 20);
3408 assert_eq!(playlist_ids(&player)[0], given[0]);
3409 assert_ne!(playlist_ids(&player), given);
3410
3411 player.process_command(PlayerCommand::SetShuffle(false));
3412 assert_eq!(playlist_ids(&player), given);
3413 }
3414
3415 #[test]
3416 fn shuffle_is_one_undo_step() {
3417 let mut player = Player::new();
3418 let ids = seed(&mut player, 10);
3419 pretend_playing(&mut player, ids[0]);
3420
3421 player.process_command(PlayerCommand::SetShuffle(true));
3422 let shuffled = playlist_ids(&player);
3423 player.process_command(PlayerCommand::Undo);
3424 assert_eq!(playlist_ids(&player), ids);
3425 assert!(!player.shared_state.play_mode().shuffle);
3426
3427 player.process_command(PlayerCommand::Redo);
3428 assert_eq!(playlist_ids(&player), shuffled);
3429 assert!(player.shared_state.play_mode().shuffle);
3430 player.process_command(PlayerCommand::SetShuffle(false));
3431 assert_eq!(
3432 playlist_ids(&player),
3433 ids,
3434 "the redone shuffle still unwinds"
3435 );
3436 }
3437
3438 #[test]
3439 fn removing_the_last_track_playing_carries_on_from_the_top_when_repeating() {
3440 let mut player = Player::new();
3441 let ids = seed(&mut player, 3);
3442 player.process_command(PlayerCommand::SetRepeat(Repeat::Queue));
3443 player.shared_state.set_cursor(Some(ids[2]));
3444
3445 player.process_command(PlayerCommand::RemoveFromPlaylist(ids[2]));
3446 assert_eq!(player.shared_state.cursor(), Some(ids[0]));
3447 }
3448
3449 #[test]
3450 fn removing_a_track_the_decoder_queued_takes_it_back() {
3451 let dir = tempfile::tempdir().unwrap();
3452 let (mut player, ids) = queued_wavs(dir.path(), &["a", "b", "c"]);
3453 assert_eq!(player.timeline.queued_after_playhead(), ids[1..]);
3454
3455 player.process_command(PlayerCommand::RemoveFromPlaylist(ids[1]));
3456 assert_eq!(player.playback_starts, 2, "restarted at the playhead");
3457 assert_eq!(player.shared_state.cursor(), Some(ids[0]));
3458 await_queued(&player);
3459 await_queued(&player);
3460 assert_eq!(player.timeline.queued_after_playhead(), vec![ids[2]]);
3461 player.process_command(PlayerCommand::Stop);
3462 }
3463
3464 #[test]
3465 fn a_track_inserted_to_play_next_is_not_skipped() {
3466 let dir = tempfile::tempdir().unwrap();
3467 let (mut player, ids) = queued_wavs(dir.path(), &["a", "b"]);
3468 let path = dir.path().join("next.wav");
3469 crate::test_utils::generate_wav(&path, 8_000, 1, 1.0, 16);
3470 let next = PlaylistItem {
3471 path,
3472 ..make_item("next")
3473 };
3474 let next_id = next.id;
3475
3476 player.process_command(PlayerCommand::InsertInPlaylist {
3477 items: vec![next],
3478 after: ids[0],
3479 });
3480 assert_eq!(player.playback_starts, 2);
3481 for _ in 0..3 {
3482 await_queued(&player);
3483 }
3484 assert_eq!(
3485 player.timeline.queued_after_playhead(),
3486 vec![next_id, ids[1]]
3487 );
3488 player.process_command(PlayerCommand::Stop);
3489 }
3490
3491 #[test]
3492 fn a_track_the_decoder_could_not_open_does_not_make_every_edit_restart() {
3493 let dir = tempfile::tempdir().unwrap();
3494 let mut player = Player::new();
3495 player.backend = Box::new(StuckBackend {
3496 rate: 8_000.0,
3497 asked: Default::default(),
3498 starts: Default::default(),
3499 });
3500 let wav = |name: &str| {
3501 let path = dir.path().join(format!("{name}.wav"));
3502 crate::test_utils::generate_wav(&path, 8_000, 1, 30.0, 16);
3503 PlaylistItem {
3504 path,
3505 ..make_item(name)
3506 }
3507 };
3508 let items = vec![wav("a"), make_item("missing"), wav("c")];
3510 let ids: Vec<_> = items.iter().map(|i| i.id).collect();
3511 player.process_command(PlayerCommand::AddToPlaylist(items));
3512 player.process_command(PlayerCommand::Play(ids[0]));
3513 await_queued(&player);
3514 await_queued(&player);
3515
3516 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("later")]));
3517 assert_eq!(
3518 player.playback_starts, 1,
3519 "nothing the decoder decided changed"
3520 );
3521 assert_eq!(player.timeline.queued_after_playhead(), vec![ids[2]]);
3522 player.process_command(PlayerCommand::Stop);
3523 }
3524
3525 #[test]
3526 fn a_track_added_after_the_decoder_reached_the_end_follows_gaplessly() {
3527 let dir = tempfile::tempdir().unwrap();
3528 let (mut player, ids) = queued_wavs(dir.path(), &["a"]);
3529 wait_for_lookahead(&player);
3530 let path = dir.path().join("b.wav");
3531 crate::test_utils::generate_wav(&path, 8_000, 1, 1.0, 16);
3532 let b = PlaylistItem {
3533 path,
3534 ..make_item("b")
3535 };
3536 let b_id = b.id;
3537
3538 player.process_command(PlayerCommand::AddToPlaylist(vec![b]));
3539 assert_eq!(player.playback_starts, 2, "restarted to queue it");
3540 assert_eq!(player.shared_state.cursor(), Some(ids[0]));
3541 await_queued(&player);
3542 await_queued(&player);
3543 assert_eq!(player.timeline.queued_after_playhead(), vec![b_id]);
3544 player.process_command(PlayerCommand::Stop);
3545 }
3546
3547 #[test]
3548 fn playing_next_a_track_still_downloading_takes_back_the_lookahead() {
3549 let dir = tempfile::tempdir().unwrap();
3550 let (mut player, ids) = queued_wavs(dir.path(), &["a", "b"]);
3551 let remote = pending_item("remote");
3552 let remote_id = remote.id;
3553
3554 player.process_command(PlayerCommand::InsertInPlaylist {
3555 items: vec![remote],
3556 after: ids[0],
3557 });
3558 assert_eq!(player.playback_starts, 2, "b is taken back out of the ring");
3559 await_queued(&player);
3560 wait_for_lookahead(&player);
3561 let session = player.session().unwrap();
3562 let step = session.lookahead.lock().last().cloned().unwrap();
3563 assert_eq!(step.next, Some(remote_id), "the decoder waits for it");
3564 assert!(player.timeline.queued_after_playhead().is_empty());
3565 player.process_command(PlayerCommand::Stop);
3566 }
3567
3568 #[test]
3569 fn gapless_playback_waits_for_a_track_still_downloading() {
3570 let dir = tempfile::tempdir().unwrap();
3571 let mut player = Player::new();
3572 player.backend = Box::new(StuckBackend {
3573 rate: 8_000.0,
3574 asked: Default::default(),
3575 starts: Default::default(),
3576 });
3577 let wav = |name: &str| {
3578 let path = dir.path().join(format!("{name}.wav"));
3579 crate::test_utils::generate_wav(&path, 8_000, 1, 1.0, 16);
3580 PlaylistItem {
3581 path,
3582 ..make_item(name)
3583 }
3584 };
3585 let items = vec![wav("a"), pending_item("arriving"), wav("c")];
3586 let ids: Vec<_> = items.iter().map(|i| i.id).collect();
3587 player.process_command(PlayerCommand::AddToPlaylist(items));
3588 player.process_command(PlayerCommand::Play(ids[0]));
3589 await_queued(&player);
3590 wait_for_lookahead(&player);
3591
3592 assert!(
3593 player.timeline.queued_after_playhead().is_empty(),
3594 "c is not queued over the track before it"
3595 );
3596 player.process_command(PlayerCommand::DecodeFinished(player.session));
3598 assert_eq!(player.shared_state.cursor(), Some(ids[1]));
3599 assert!(player.waiting().is_some());
3600 player.process_command(PlayerCommand::Stop);
3601 }
3602
3603 fn wait_for_lookahead(player: &Player) {
3605 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
3606 while player
3607 .session()
3608 .is_some_and(|s| s.lookahead.lock().is_empty())
3609 {
3610 assert!(
3611 std::time::Instant::now() < deadline,
3612 "the decoder never looked ahead"
3613 );
3614 std::thread::yield_now();
3615 }
3616 }
3617
3618 #[test]
3619 fn an_edit_after_what_the_decoder_queued_leaves_playback_alone() {
3620 let dir = tempfile::tempdir().unwrap();
3621 let (mut player, ids) = wavs_playing(dir.path(), &["a", "b"], 30.0);
3622 await_queued(&player);
3623 await_queued(&player);
3624 let path = dir.path().join("later.wav");
3625 crate::test_utils::generate_wav(&path, 8_000, 1, 1.0, 16);
3626
3627 player.process_command(PlayerCommand::AddToPlaylist(vec![PlaylistItem {
3628 path,
3629 ..make_item("later")
3630 }]));
3631 assert_eq!(player.playback_starts, 1);
3632 assert_eq!(player.timeline.queued_after_playhead(), vec![ids[1]]);
3633 player.process_command(PlayerCommand::Stop);
3634 }
3635
3636 #[test]
3637 fn undoing_the_add_of_a_track_on_its_way_forgets_it() {
3638 let dir = tempfile::tempdir().unwrap();
3639 let (mut player, _, starts) = downloading_wav(dir.path());
3640 let path = dir.path().join("added.wav");
3641 crate::test_utils::generate_wav(&path, 8_000, 1, 1.0, 16);
3642 let added = PlaylistItem {
3643 path,
3644 state: ItemState::Pending,
3645 ..make_item("added")
3646 };
3647 let added_id = added.id;
3648 player.process_command(PlayerCommand::AddToPlaylist(vec![added]));
3649 player.process_command(PlayerCommand::Play(added_id));
3650
3651 player.process_command(PlayerCommand::Undo);
3652 assert!(player.waiting().is_none());
3653 assert_eq!(player.shared_state.playback_state(), PlaybackState::Stopped);
3654
3655 player.process_command(PlayerCommand::Redo);
3657 player
3658 .shared_state
3659 .update_item_state(added_id, ItemState::Ready);
3660 player.process_command(PlayerCommand::TrackReady(added_id));
3661 assert!(player.session().is_none());
3662 assert_eq!(starts.load(Ordering::Relaxed), 0);
3663 }
3664
3665 pub(super) struct Rng(pub(super) u64);
3668
3669 impl Rng {
3670 pub(super) fn next(&mut self) -> u64 {
3671 self.0 ^= self.0 << 13;
3672 self.0 ^= self.0 >> 7;
3673 self.0 ^= self.0 << 17;
3674 self.0
3675 }
3676
3677 pub(super) fn below(&mut self, n: usize) -> usize {
3678 (self.next() % n.max(1) as u64) as usize
3679 }
3680
3681 pub(super) fn coin(&mut self) -> bool {
3682 self.next() & 1 == 1
3683 }
3684 }
3685
3686 pub(super) fn asks_to_play(cmd: &PlayerCommand) -> bool {
3689 cmd.asks_to_play()
3690 }
3691
3692 pub(super) fn check_invariants(
3695 player: &Player,
3696 wanted_before: bool,
3697 asked: bool,
3698 ) -> Result<(), String> {
3699 let state = &player.shared_state;
3700 let ids = playlist_ids(player);
3701 let listed = |id: QueueItemId| ids.contains(&id);
3702 let mut seen = std::collections::HashSet::new();
3703 if !ids.iter().all(|id| seen.insert(*id)) {
3704 return Err("an item is in the playlist twice".into());
3705 }
3706 let cursor = state.cursor();
3707 if cursor.is_some_and(|c| !listed(c)) {
3708 return Err("the cursor is on an item not in the playlist".into());
3709 }
3710 let info = state.track_info();
3711 match (&player.transport, &info) {
3712 (Transport::Loaded(_), None) => return Err("loaded, but no track published".into()),
3713 (Transport::Loaded(_), Some(info)) => {
3714 if !listed(info.id) {
3715 return Err("playing an item not in the playlist".into());
3716 }
3717 if cursor != Some(info.id) {
3718 return Err("the cursor is not on what is playing".into());
3719 }
3720 if !matches!(
3721 state.playback_state(),
3722 PlaybackState::Playing | PlaybackState::Paused
3723 ) {
3724 return Err("loaded, but published as stopped".into());
3725 }
3726 }
3727 (_, Some(_)) => return Err("a track is published with nothing loaded".into()),
3728 (Transport::Waiting(w), None) => {
3729 if cursor != Some(w.id) || !listed(w.id) {
3730 return Err("waiting on an item that is not the cursor's".into());
3731 }
3732 let expected = match w.start {
3733 Run::Playing => PlaybackState::Stopped,
3734 Run::Paused => PlaybackState::Paused,
3735 };
3736 if state.playback_state() != expected {
3737 return Err("a wait published in the wrong state".into());
3738 }
3739 }
3740 (Transport::Idle, None) => {
3741 if state.playback_state() != PlaybackState::Stopped {
3742 return Err("idle, but not published as stopped".into());
3743 }
3744 }
3745 }
3746 if state.is_waiting() != player.waiting().is_some() {
3747 return Err("the published wait disagrees with the player".into());
3748 }
3749 if player
3750 .timeline
3751 .queued_after_playhead()
3752 .into_iter()
3753 .any(|id| !listed(id))
3754 {
3755 return Err("the decoder has queued an item no longer in the playlist".into());
3756 }
3757 if !asked && !wanted_before && state.wants_to_play() {
3758 return Err("started without being asked".into());
3759 }
3760 if state.play_mode() != player.mode {
3761 return Err("the published mode disagrees with the player".into());
3762 }
3763 if !player.mode.shuffle && state.shuffle_order().iter().any(|(_, pre)| pre.is_some()) {
3764 return Err("shuffle is off, but an item remembers a place to go back to".into());
3765 }
3766 if player.mode.repeat == Repeat::Off
3767 && let Some(session) = player.session()
3768 && let Some(playhead) = player.timeline.playhead()
3769 && session.track.id == playhead.id
3770 && session
3771 .lookahead
3772 .lock()
3773 .iter()
3774 .filter(|step| step.boundary > playhead.boundary)
3775 .any(|step| step.wrapped || step.next == Some(step.after))
3776 {
3777 return Err("repeat is off, but the decoder has queued a wrap".into());
3778 }
3779 Ok(())
3780 }
3781
3782 #[test]
3783 fn random_use_keeps_the_player_honest() {
3784 let dir = tempfile::tempdir().unwrap();
3785 let path = dir.path().join("t.wav");
3786 crate::test_utils::generate_wav(&path, 8_000, 1, 0.5, 16);
3787
3788 for seed in 1..=40u64 {
3789 let mut rng = Rng(seed.wrapping_mul(0x9E37_79B9_7F4A_7C15) | 1);
3790 let mut player = Player::new();
3791 player.backend = Box::new(StuckBackend {
3792 rate: 8_000.0,
3793 asked: Default::default(),
3794 starts: Default::default(),
3795 });
3796 let fresh = |rng: &mut Rng| {
3797 let state = match rng.below(3) {
3798 0 => ItemState::Pending,
3799 _ => ItemState::Ready,
3800 };
3801 PlaylistItem {
3802 path: path.clone(),
3803 state,
3804 ..make_item("t")
3805 }
3806 };
3807 let start: Vec<_> = (0..5).map(|_| fresh(&mut rng)).collect();
3808 player.process_command(PlayerCommand::AddToPlaylist(start));
3809
3810 let mut history = Vec::new();
3811 for step in 0..150 {
3812 let ids = playlist_ids(&player);
3813 let pick = |rng: &mut Rng| ids.get(rng.below(ids.len())).copied();
3814 let cmd = match rng.below(23) {
3815 0 => pick(&mut rng).map(PlayerCommand::Play),
3816 1 => pick(&mut rng).map(|id| PlayerCommand::Cue {
3817 id,
3818 position_ms: if rng.coin() { 0 } else { 200 },
3819 play: rng.coin(),
3820 }),
3821 2 => Some(PlayerCommand::Pause),
3822 3 => Some(PlayerCommand::Resume),
3823 4 => Some(PlayerCommand::NextTrack),
3824 5 => Some(PlayerCommand::PrevTrack),
3825 6 => Some(PlayerCommand::Seek(100)),
3826 7 => pick(&mut rng).map(PlayerCommand::RemoveFromPlaylist),
3827 8 => Some(PlayerCommand::RemoveFromPlaylistBatch(
3828 (0..2).filter_map(|_| pick(&mut rng)).collect(),
3829 )),
3830 9 => pick(&mut rng).zip(pick(&mut rng)).map(|(id, target)| {
3831 PlayerCommand::MoveInPlaylist {
3832 id,
3833 target,
3834 after: rng.coin(),
3835 }
3836 }),
3837 10 => pick(&mut rng).map(|after| PlayerCommand::InsertInPlaylist {
3838 items: vec![fresh(&mut rng)],
3839 after,
3840 }),
3841 11 => Some(PlayerCommand::AddToPlaylist(vec![fresh(&mut rng)])),
3842 12 => Some(PlayerCommand::Undo),
3843 13 => Some(PlayerCommand::Redo),
3844 14 | 15 => {
3845 let (items, _) = player.shared_state.snapshot_playlist();
3846 let pending: Vec<_> = items
3847 .iter()
3848 .filter(|i| matches!(i.state, ItemState::Pending))
3849 .map(|i| i.id)
3850 .collect();
3851 pending.get(rng.below(pending.len())).map(|&id| {
3852 if rng.below(4) == 0 {
3853 player
3854 .shared_state
3855 .update_item_state(id, ItemState::Failed("gone".into()));
3856 PlayerCommand::TrackFailed(id)
3857 } else {
3858 player.shared_state.update_item_state(id, ItemState::Ready);
3859 PlayerCommand::TrackReady(id)
3860 }
3861 })
3862 }
3863 16 => Some(PlayerCommand::DecodeFinished(player.session)),
3864 17 => Some(PlayerCommand::DecodeFinished(
3865 player.session.wrapping_sub(1),
3866 )),
3867 18 => Some(PlayerCommand::ClearPlaylist),
3868 20 => Some(PlayerCommand::SetShuffle(rng.coin())),
3869 21 => Some(PlayerCommand::SetRepeat(
3870 [Repeat::Off, Repeat::Queue, Repeat::One][rng.below(3)],
3871 )),
3872 _ => Some(PlayerCommand::ReplacePlaylist {
3873 items: (0..3).map(|_| fresh(&mut rng)).collect(),
3874 start: rng.below(4),
3875 position_ms: if rng.coin() { 0 } else { 200 },
3876 play: rng.coin(),
3877 }),
3878 };
3879 let Some(cmd) = cmd else { continue };
3880 let label = format!("{cmd:?}");
3881 let asked = asks_to_play(&cmd);
3882 let replaced_shuffled = player.mode.shuffle
3883 && matches!(&cmd, PlayerCommand::ReplacePlaylist { items, .. } if items.len() > 1);
3884 let wanted_before = player.shared_state.wants_to_play();
3885 player.process_command(cmd);
3886 while let Ok(sent) = player.commands.rx.try_recv() {
3890 player.process_command(sent);
3891 }
3892 player.update_playback_state();
3893 history.push(label);
3894 let broken = check_invariants(&player, wanted_before, asked)
3895 .err()
3896 .or_else(|| {
3897 (replaced_shuffled
3898 && player
3899 .shared_state
3900 .shuffle_order()
3901 .iter()
3902 .all(|(_, pre)| pre.is_none()))
3903 .then(|| "a queue replaced while shuffled plays in order".to_string())
3904 });
3905 if let Some(broken) = broken {
3906 let tail = history[history.len().saturating_sub(8)..].join("\n ");
3907 panic!("seed {seed}, step {step}: {broken}\nlast commands:\n {tail}");
3908 }
3909 }
3910 player.process_command(PlayerCommand::Stop);
3911 }
3912 }
3913
3914 #[test]
3915 fn a_pause_reports_where_the_fade_went_silent() {
3916 use std::sync::atomic::AtomicBool;
3917
3918 struct FadingEngine {
3919 running: Arc<AtomicBool>,
3920 silent: Arc<AtomicBool>,
3921 }
3922 impl AudioEngineHandle for FadingEngine {
3923 fn start(&self) -> Result<(), BackendError> {
3924 self.running.store(true, Ordering::Relaxed);
3925 Ok(())
3926 }
3927 fn stop(&self) -> Result<(), BackendError> {
3928 self.running.store(false, Ordering::Relaxed);
3929 Ok(())
3930 }
3931 fn is_running(&self) -> bool {
3932 self.running.load(Ordering::Relaxed)
3933 }
3934 fn fade_out(&self) {}
3935 fn fade_in(&self) -> Result<(), BackendError> {
3936 Ok(())
3937 }
3938 fn is_silent(&self) -> bool {
3939 self.silent.load(Ordering::Relaxed)
3940 }
3941 }
3942
3943 let running = Arc::new(AtomicBool::new(true));
3944 let silent = Arc::new(AtomicBool::new(false));
3945 let mut player = Player::new();
3946 player.transport = Transport::Loaded(test_session(
3947 QueueItemId::new(),
3948 Box::new(FadingEngine {
3949 running: running.clone(),
3950 silent: silent.clone(),
3951 }),
3952 ));
3953 player
3954 .shared_state
3955 .set_playback_state(PlaybackState::Playing);
3956 player.shared_state.set_position_ms(5_000);
3957
3958 let (reply, answer) = crossbeam_channel::bounded(1);
3959 player.process_command(PlayerCommand::PauseAndReport(reply));
3960
3961 if crate::config::Config::cached().playback.fade_on_pause {
3962 assert!(answer.try_recv().is_err(), "not while the fade is audible");
3963 player.shared_state.set_position_ms(5_150);
3965 player.update_playback_state();
3966 assert!(answer.try_recv().is_err());
3967 silent.store(true, Ordering::Relaxed);
3968 player.update_playback_state();
3969 assert!(!running.load(Ordering::Relaxed));
3970 assert_eq!(answer.try_recv().unwrap(), 5_150);
3971 } else {
3972 assert_eq!(answer.try_recv().unwrap(), 5_000);
3973 }
3974 }
3975
3976 #[test]
3977 fn the_server_hears_each_turn_playback_takes() {
3978 use PlaybackReportState::{Paused, Playing, Stopped};
3979 use history::PlaybackReport;
3980
3981 let dir = tempfile::tempdir().unwrap();
3982 let path = dir.path().join("t.wav");
3983 crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
3984
3985 let mut player = Player::new();
3986 player.backend = Box::new(StuckBackend {
3987 rate: 8_000.0,
3988 asked: Default::default(),
3989 starts: Default::default(),
3990 });
3991 let (recorder, events) = PlayRecorder::capture();
3992 player.history = Some(recorder);
3993
3994 let item = PlaylistItem {
3995 db_id: Some(5),
3996 path,
3997 ..make_item("t")
3998 };
3999 let id = item.id;
4000 player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
4001 player.process_command(PlayerCommand::Play(id));
4002 player.process_command(PlayerCommand::Pause);
4003 player.process_command(PlayerCommand::Seek(4_000));
4004 player.process_command(PlayerCommand::Resume);
4005 player.process_command(PlayerCommand::Stop);
4006
4007 let report = |state, position_ms| {
4008 PlayEvent::Playback(PlaybackReport {
4009 track_id: 5,
4010 state,
4011 position_ms,
4012 })
4013 };
4014 assert_eq!(
4015 events.try_iter().collect::<Vec<_>>(),
4016 vec![
4017 PlayEvent::Started {
4018 track_id: 5,
4019 position_ms: 0
4020 },
4021 report(Paused, 0),
4022 report(Paused, 4_000),
4023 report(Playing, 4_000),
4024 report(Stopped, 4_000),
4025 PlayEvent::Finished {
4026 track_id: 5,
4027 listened_ms: 0
4028 },
4029 ]
4030 );
4031 }
4032
4033 #[test]
4034 fn removing_the_playing_track_resumes_at_its_successor() {
4035 let mut player = Player::new();
4036 let ids = seed(&mut player, 5);
4037 pretend_playing(&mut player, ids[2]);
4038
4039 player.process_command(PlayerCommand::RemoveFromPlaylist(ids[2]));
4040
4041 assert_eq!(
4042 player.shared_state.cursor(),
4043 Some(ids[3]),
4044 "playback must continue at the next track, not restart the queue"
4045 );
4046 assert_eq!(player.playback_starts, 1);
4047 }
4048
4049 #[test]
4050 fn removing_the_paused_track_moves_on_paused() {
4051 let mut player = Player::new();
4052 let ids = seed(&mut player, 3);
4053 pretend_playing(&mut player, ids[1]);
4054 if let Transport::Loaded(session) = &mut player.transport {
4055 session.run = Run::Paused;
4056 }
4057
4058 player.process_command(PlayerCommand::RemoveFromPlaylist(ids[1]));
4059 assert_eq!(player.shared_state.cursor(), Some(ids[2]));
4060 assert!(!player.shared_state.wants_to_play(), "still paused");
4061 }
4062
4063 #[test]
4064 fn removing_the_cursor_with_nothing_loaded_starts_nothing() {
4065 let mut player = Player::new();
4066 let ids = seed(&mut player, 3);
4067 player.shared_state.set_cursor(Some(ids[1]));
4068
4069 player.process_command(PlayerCommand::RemoveFromPlaylist(ids[1]));
4070 assert_eq!(player.shared_state.cursor(), Some(ids[2]));
4071 assert_eq!(player.playback_starts, 0);
4072 }
4073
4074 #[test]
4075 fn removing_the_first_playing_track_resumes_at_the_new_first() {
4076 let mut player = Player::new();
4077 let ids = seed(&mut player, 3);
4078 pretend_playing(&mut player, ids[0]);
4079
4080 player.process_command(PlayerCommand::RemoveFromPlaylist(ids[0]));
4081
4082 assert_eq!(player.shared_state.cursor(), Some(ids[1]));
4083 }
4084
4085 #[test]
4086 fn next_track_parks_on_a_track_that_has_not_downloaded_yet() {
4087 let mut player = Player::new();
4088 let playing = make_item("playing");
4089 let waiting = pending_item("waiting");
4090 let later = make_item("later");
4091 let (playing_id, waiting_id) = (playing.id, waiting.id);
4092 player.process_command(PlayerCommand::AddToPlaylist(vec![playing, waiting, later]));
4093 pretend_playing(&mut player, playing_id);
4094
4095 player.process_command(PlayerCommand::DecodeFinished(player.session));
4096
4097 assert_eq!(
4098 player.shared_state.cursor(),
4099 Some(waiting_id),
4100 "the cursor parks on the track being fetched"
4101 );
4102 assert_eq!(
4103 player.playback_starts, 0,
4104 "nothing to play until its bytes land"
4105 );
4106
4107 player
4110 .shared_state
4111 .update_item_state(waiting_id, ItemState::Ready);
4112 player.process_command(PlayerCommand::TrackReady(waiting_id));
4113
4114 assert_eq!(player.playback_starts, 1);
4115 assert_eq!(player.shared_state.cursor(), Some(waiting_id));
4116 }
4117
4118 #[test]
4119 fn a_download_that_cannot_land_moves_the_cursor_on() {
4120 let mut player = Player::new();
4121 let waiting = pending_item("waiting");
4122 let later = make_item("later");
4123 let (waiting_id, later_id) = (waiting.id, later.id);
4124 player.process_command(PlayerCommand::AddToPlaylist(vec![waiting, later]));
4125
4126 player.process_command(PlayerCommand::Play(waiting_id));
4127 assert_eq!(player.playback_starts, 0, "nothing to play yet");
4128
4129 player
4131 .shared_state
4132 .update_item_state(waiting_id, ItemState::Failed("remote unavailable".into()));
4133 player.process_command(PlayerCommand::TrackFailed(waiting_id));
4134
4135 assert_eq!(
4136 player.shared_state.cursor(),
4137 Some(later_id),
4138 "the queue moves past a track that can never load"
4139 );
4140 assert_eq!(player.playback_starts, 1);
4141 }
4142
4143 #[test]
4144 fn a_queue_that_can_never_load_stops_rather_than_waiting() {
4145 let mut player = Player::new();
4146 let first = pending_item("first");
4147 let second = pending_item("second");
4148 let (first_id, second_id) = (first.id, second.id);
4149 player.process_command(PlayerCommand::AddToPlaylist(vec![first, second]));
4150
4151 player.process_command(PlayerCommand::Play(first_id));
4152 for id in [first_id, second_id] {
4153 player
4154 .shared_state
4155 .update_item_state(id, ItemState::Failed("remote unavailable".into()));
4156 player.process_command(PlayerCommand::TrackFailed(id));
4157 }
4158
4159 assert_eq!(player.playback_starts, 0);
4160 assert_eq!(
4161 player.shared_state.playback_state(),
4162 PlaybackState::Stopped,
4163 "a stop the UI can see, not an indefinite wait for TrackReady"
4164 );
4165 }
4166
4167 #[test]
4168 fn a_failure_elsewhere_in_the_queue_leaves_the_cursor_alone() {
4169 let mut player = Player::new();
4170 let waiting = pending_item("waiting");
4171 let other = pending_item("other");
4172 let (waiting_id, other_id) = (waiting.id, other.id);
4173 player.process_command(PlayerCommand::AddToPlaylist(vec![waiting, other]));
4174 player.process_command(PlayerCommand::Play(waiting_id));
4175
4176 player
4177 .shared_state
4178 .update_item_state(other_id, ItemState::Failed("remote unavailable".into()));
4179 player.process_command(PlayerCommand::TrackFailed(other_id));
4180
4181 assert_eq!(
4182 player.shared_state.cursor(),
4183 Some(waiting_id),
4184 "a track still downloading keeps the cursor"
4185 );
4186 }
4187
4188 #[test]
4189 fn batch_delete_containing_the_cursor_restarts_the_engine_once() {
4190 let mut player = Player::new();
4191 let ids = seed(&mut player, 5);
4192 pretend_playing(&mut player, ids[2]);
4193
4194 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![
4195 ids[1], ids[2], ids[3],
4196 ]));
4197
4198 assert_eq!(playlist_titles(&player), vec!["t0", "t4"]);
4199 assert_eq!(player.shared_state.cursor(), Some(ids[4]));
4200 assert_eq!(
4201 player.playback_starts, 1,
4202 "one resume for the whole selection, not one per deleted track"
4203 );
4204 }
4205
4206 #[test]
4207 fn batch_delete_below_the_cursor_leaves_playback_alone() {
4208 let mut player = Player::new();
4209 let ids = seed(&mut player, 4);
4210 player.shared_state.set_cursor(Some(ids[0]));
4211
4212 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![ids[2], ids[3]]));
4213
4214 assert_eq!(player.shared_state.cursor(), Some(ids[0]));
4215 assert_eq!(player.playback_starts, 0);
4216 }
4217
4218 #[test]
4219 fn undo_of_a_batch_delete_restores_the_original_order() {
4220 let mut player = Player::new();
4224 let items = vec![
4225 make_item("A"),
4226 make_item("B"),
4227 make_item("C"),
4228 make_item("D"),
4229 ];
4230 let (b_id, c_id) = (items[1].id, items[2].id);
4231 player.process_command(PlayerCommand::AddToPlaylist(items));
4232
4233 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![c_id, b_id]));
4234 assert_eq!(playlist_titles(&player), vec!["A", "D"]);
4235
4236 player.process_command(PlayerCommand::Undo);
4237 assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
4238 }
4239
4240 #[test]
4243 fn undo_add_removes_items() {
4244 let mut player = Player::new();
4245 let items = vec![make_item("A"), make_item("B")];
4246 let ids: Vec<_> = items.iter().map(|i| i.id).collect();
4247
4248 player.process_command(PlayerCommand::AddToPlaylist(items));
4249 assert_eq!(playlist_ids(&player), ids);
4250 assert!(player.undo_stack().can_undo());
4251
4252 player.process_command(PlayerCommand::Undo);
4253 assert!(playlist_ids(&player).is_empty());
4254 assert!(player.undo_stack().can_redo());
4255 }
4256
4257 #[test]
4258 fn redo_add_restores_items() {
4259 let mut player = Player::new();
4260 let items = vec![make_item("A"), make_item("B")];
4261
4262 player.process_command(PlayerCommand::AddToPlaylist(items));
4263 player.process_command(PlayerCommand::Undo);
4264 assert!(playlist_ids(&player).is_empty());
4265
4266 player.process_command(PlayerCommand::Redo);
4267 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
4268 }
4269
4270 #[test]
4273 fn undo_remove_restores_item_at_position() {
4274 let mut player = Player::new();
4275 let items = vec![make_item("A"), make_item("B"), make_item("C")];
4276 let b_id = items[1].id;
4277
4278 player.process_command(PlayerCommand::AddToPlaylist(items));
4279 player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
4280 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
4281
4282 player.process_command(PlayerCommand::Undo);
4283 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4284 }
4285
4286 #[test]
4287 fn undo_remove_first_item() {
4288 let mut player = Player::new();
4289 let items = vec![make_item("A"), make_item("B")];
4290 let a_id = items[0].id;
4291
4292 player.process_command(PlayerCommand::AddToPlaylist(items));
4293 player.process_command(PlayerCommand::RemoveFromPlaylist(a_id));
4294 assert_eq!(playlist_titles(&player), vec!["B"]);
4295
4296 player.process_command(PlayerCommand::Undo);
4297 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
4298 }
4299
4300 #[test]
4301 fn undo_batch_remove_restores_all() {
4302 let mut player = Player::new();
4303 let items = vec![
4304 make_item("A"),
4305 make_item("B"),
4306 make_item("C"),
4307 make_item("D"),
4308 ];
4309 let b_id = items[1].id;
4310 let c_id = items[2].id;
4311
4312 player.process_command(PlayerCommand::AddToPlaylist(items));
4313 let version_before = player.shared_state.playlist_version();
4314 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![b_id, c_id]));
4315 assert_eq!(playlist_titles(&player), vec!["A", "D"]);
4316 assert_eq!(
4319 player.shared_state.playlist_version(),
4320 version_before + 1,
4321 "batch removal must bump the playlist version exactly once"
4322 );
4323
4324 player.process_command(PlayerCommand::Undo);
4326 assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
4327 }
4328
4329 #[test]
4330 fn redo_batch_remove() {
4331 let mut player = Player::new();
4332 let items = vec![make_item("A"), make_item("B"), make_item("C")];
4333 let a_id = items[0].id;
4334 let b_id = items[1].id;
4335
4336 player.process_command(PlayerCommand::AddToPlaylist(items));
4337 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![a_id, b_id]));
4338 player.process_command(PlayerCommand::Undo);
4339 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4340
4341 player.process_command(PlayerCommand::Redo);
4342 assert_eq!(playlist_titles(&player), vec!["C"]);
4343 }
4344
4345 #[test]
4346 fn redo_remove() {
4347 let mut player = Player::new();
4348 let items = vec![make_item("A"), make_item("B"), make_item("C")];
4349 let b_id = items[1].id;
4350
4351 player.process_command(PlayerCommand::AddToPlaylist(items));
4352 player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
4353 player.process_command(PlayerCommand::Undo);
4354 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4355
4356 player.process_command(PlayerCommand::Redo);
4357 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
4358 }
4359
4360 #[test]
4363 fn undo_insert_removes_inserted_items() {
4364 let mut player = Player::new();
4365 let items = vec![make_item("A"), make_item("C")];
4366 let a_id = items[0].id;
4367
4368 player.process_command(PlayerCommand::AddToPlaylist(items));
4369
4370 let inserted = vec![make_item("B")];
4371 player.process_command(PlayerCommand::InsertInPlaylist {
4372 items: inserted,
4373 after: a_id,
4374 });
4375 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4376
4377 player.process_command(PlayerCommand::Undo);
4378 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
4379 }
4380
4381 #[test]
4384 fn undo_move_restores_position() {
4385 let mut player = Player::new();
4386 let items = vec![make_item("A"), make_item("B"), make_item("C")];
4387 let a_id = items[0].id;
4388 let c_id = items[2].id;
4389
4390 player.process_command(PlayerCommand::AddToPlaylist(items));
4391
4392 player.process_command(PlayerCommand::MoveInPlaylist {
4394 id: a_id,
4395 target: c_id,
4396 after: true,
4397 });
4398 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
4399
4400 player.process_command(PlayerCommand::Undo);
4401 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4402 }
4403
4404 #[test]
4405 fn redo_move() {
4406 let mut player = Player::new();
4407 let items = vec![make_item("A"), make_item("B"), make_item("C")];
4408 let a_id = items[0].id;
4409 let c_id = items[2].id;
4410
4411 player.process_command(PlayerCommand::AddToPlaylist(items));
4412 player.process_command(PlayerCommand::MoveInPlaylist {
4413 id: a_id,
4414 target: c_id,
4415 after: true,
4416 });
4417 player.process_command(PlayerCommand::Undo);
4418 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4419
4420 player.process_command(PlayerCommand::Redo);
4421 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
4422 }
4423
4424 #[test]
4427 fn undo_batch_move() {
4428 let mut player = Player::new();
4429 let items = vec![
4430 make_item("A"),
4431 make_item("B"),
4432 make_item("C"),
4433 make_item("D"),
4434 ];
4435 let a_id = items[0].id;
4436 let b_id = items[1].id;
4437 let d_id = items[3].id;
4438
4439 player.process_command(PlayerCommand::AddToPlaylist(items));
4440
4441 player.process_command(PlayerCommand::MoveItemsInPlaylist {
4443 ids: vec![a_id, b_id],
4444 target: d_id,
4445 after: true,
4446 });
4447 assert_eq!(playlist_titles(&player), vec!["C", "D", "A", "B"]);
4448
4449 player.process_command(PlayerCommand::Undo);
4450 assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
4451 }
4452
4453 #[test]
4456 fn undo_clear_restores_playlist() {
4457 let mut player = Player::new();
4458 let items = vec![make_item("A"), make_item("B"), make_item("C")];
4459
4460 player.process_command(PlayerCommand::AddToPlaylist(items));
4461 player.process_command(PlayerCommand::ClearPlaylist);
4462 assert!(playlist_ids(&player).is_empty());
4463
4464 player.process_command(PlayerCommand::Undo);
4465 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4466 }
4467
4468 #[test]
4473 fn undoing_a_replace_does_not_leave_the_engine_on_an_orphaned_track() {
4474 let mut player = Player::new();
4475 let original = seed(&mut player, 3);
4476 player.shared_state.set_cursor(Some(original[0]));
4477 pretend_playing(&mut player, original[0]);
4478
4479 let replacement = vec![make_item("something else")];
4480 let orphan = replacement[0].id;
4481 player.process_command(PlayerCommand::ReplacePlaylist {
4482 items: replacement,
4483 start: 0,
4484 position_ms: 0,
4485 play: true,
4486 });
4487 pretend_playing(&mut player, orphan);
4489
4490 player.process_command(PlayerCommand::Undo);
4491
4492 assert_eq!(playlist_ids(&player), original, "the queue comes back");
4493 assert!(
4494 player.shared_state.get_item(orphan).is_none(),
4495 "and the replacement is gone from it"
4496 );
4497 assert!(
4498 playing_id(&player).is_none_or(|id| player.shared_state.get_item(id).is_some()),
4499 "so nothing may still be playing out of it"
4500 );
4501 }
4502
4503 #[test]
4505 fn undoing_an_add_does_not_leave_the_engine_on_a_removed_track() {
4506 let mut player = Player::new();
4507 seed(&mut player, 2);
4508 let added = seed(&mut player, 1);
4509 pretend_playing(&mut player, added[0]);
4510
4511 player.process_command(PlayerCommand::Undo);
4512
4513 assert!(player.shared_state.get_item(added[0]).is_none());
4514 assert!(
4515 playing_id(&player).is_none_or(|id| player.shared_state.get_item(id).is_some()),
4516 "the engine cannot be left on the item the undo removed"
4517 );
4518 }
4519
4520 #[test]
4522 fn undoing_a_move_leaves_playback_alone() {
4523 let mut player = Player::new();
4524 let ids = seed(&mut player, 3);
4525 player.shared_state.set_cursor(Some(ids[0]));
4526 pretend_playing(&mut player, ids[0]);
4527 let starts = player.playback_starts;
4528
4529 player.process_command(PlayerCommand::MoveInPlaylist {
4530 id: ids[2],
4531 target: ids[0],
4532 after: false,
4533 });
4534 player.process_command(PlayerCommand::Undo);
4535
4536 assert_eq!(playlist_ids(&player), ids);
4537 assert_eq!(playing_id(&player), Some(ids[0]), "still on the same track");
4538 assert_eq!(player.playback_starts, starts, "and not restarted");
4539 }
4540
4541 #[test]
4542 fn redo_clear() {
4543 let mut player = Player::new();
4544 let items = vec![make_item("A"), make_item("B")];
4545
4546 player.process_command(PlayerCommand::AddToPlaylist(items));
4547 player.process_command(PlayerCommand::ClearPlaylist);
4548 player.process_command(PlayerCommand::Undo);
4549 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
4550
4551 player.process_command(PlayerCommand::Redo);
4552 assert!(playlist_ids(&player).is_empty());
4553 }
4554
4555 #[test]
4558 fn multiple_undos_in_sequence() {
4559 let mut player = Player::new();
4560
4561 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("A")]));
4562 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
4563 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("C")]));
4564 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4565
4566 player.process_command(PlayerCommand::Undo);
4567 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
4568
4569 player.process_command(PlayerCommand::Undo);
4570 assert_eq!(playlist_titles(&player), vec!["A"]);
4571
4572 player.process_command(PlayerCommand::Undo);
4573 assert!(playlist_ids(&player).is_empty());
4574 }
4575
4576 #[test]
4577 fn undo_redo_undo_cycle() {
4578 let mut player = Player::new();
4579 let items = vec![make_item("A"), make_item("B")];
4580
4581 player.process_command(PlayerCommand::AddToPlaylist(items));
4582 player.process_command(PlayerCommand::Undo);
4583 assert!(playlist_ids(&player).is_empty());
4584
4585 player.process_command(PlayerCommand::Redo);
4586 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
4587
4588 player.process_command(PlayerCommand::Undo);
4589 assert!(playlist_ids(&player).is_empty());
4590 }
4591
4592 #[test]
4593 fn new_action_clears_redo_stack() {
4594 let mut player = Player::new();
4595 let items = vec![make_item("A")];
4596
4597 player.process_command(PlayerCommand::AddToPlaylist(items));
4598 player.process_command(PlayerCommand::Undo);
4599 assert!(player.undo_stack().can_redo());
4600
4601 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
4603 assert!(!player.undo_stack().can_redo());
4604 }
4605
4606 #[test]
4607 fn undo_on_empty_stack_is_noop() {
4608 let mut player = Player::new();
4609 player.process_command(PlayerCommand::Undo);
4610 assert!(playlist_ids(&player).is_empty());
4611 }
4612
4613 #[test]
4614 fn redo_on_empty_stack_is_noop() {
4615 let mut player = Player::new();
4616 player.process_command(PlayerCommand::Redo);
4617 assert!(playlist_ids(&player).is_empty());
4618 }
4619
4620 #[test]
4623 fn playback_commands_not_undoable() {
4624 let mut player = Player::new();
4625 player.process_command(PlayerCommand::Pause);
4626 player.process_command(PlayerCommand::Resume);
4627 player.process_command(PlayerCommand::NextTrack);
4628 player.process_command(PlayerCommand::PrevTrack);
4629 assert!(!player.undo_stack().can_undo());
4630 }
4631
4632 #[test]
4633 fn update_paths_not_undoable() {
4634 let mut player = Player::new();
4635 let items = vec![make_item("A")];
4636 let id = items[0].id;
4637 player.process_command(PlayerCommand::AddToPlaylist(items));
4638
4639 let undo_count = player.undo_stack().undo_len();
4640 player.process_command(PlayerCommand::UpdatePaths(vec![(
4641 id,
4642 PathBuf::from("/new/path.flac"),
4643 )]));
4644 assert_eq!(player.undo_stack().undo_len(), undo_count);
4645 }
4646
4647 #[test]
4650 fn add_remove_undo_undo_produces_original() {
4651 let mut player = Player::new();
4652 let items = vec![make_item("A"), make_item("B"), make_item("C")];
4653 let b_id = items[1].id;
4654 let original_titles = vec!["A", "B", "C"];
4655
4656 player.process_command(PlayerCommand::AddToPlaylist(items));
4657 player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
4658 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
4659
4660 player.process_command(PlayerCommand::Undo);
4662 assert_eq!(playlist_titles(&player), original_titles);
4663
4664 player.process_command(PlayerCommand::Undo);
4666 assert!(playlist_ids(&player).is_empty());
4667 }
4668
4669 #[test]
4670 fn interleaved_adds_and_moves_undo() {
4671 let mut player = Player::new();
4672 let items = vec![make_item("A"), make_item("B"), make_item("C")];
4673 let a_id = items[0].id;
4674 let c_id = items[2].id;
4675
4676 player.process_command(PlayerCommand::AddToPlaylist(items));
4677
4678 player.process_command(PlayerCommand::MoveInPlaylist {
4680 id: a_id,
4681 target: c_id,
4682 after: true,
4683 });
4684 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
4685
4686 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("D")]));
4688 assert_eq!(playlist_titles(&player), vec!["B", "C", "A", "D"]);
4689
4690 player.process_command(PlayerCommand::Undo);
4692 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
4693
4694 player.process_command(PlayerCommand::Undo);
4696 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
4697 }
4698
4699 #[test]
4704 fn stop_engine_drops_engine_synchronously() {
4705 use std::sync::atomic::{AtomicBool, Ordering};
4706
4707 struct MockEngine {
4708 dropped: Arc<AtomicBool>,
4709 }
4710 impl AudioEngineHandle for MockEngine {
4711 fn start(&self) -> Result<(), BackendError> {
4712 Ok(())
4713 }
4714 fn stop(&self) -> Result<(), BackendError> {
4715 Ok(())
4716 }
4717 fn is_running(&self) -> bool {
4718 false
4719 }
4720 fn fade_out(&self) {}
4721 fn fade_in(&self) -> Result<(), BackendError> {
4722 Ok(())
4723 }
4724 fn is_silent(&self) -> bool {
4725 false
4726 }
4727 }
4728 impl Drop for MockEngine {
4729 fn drop(&mut self) {
4730 self.dropped.store(true, Ordering::SeqCst);
4731 }
4732 }
4733
4734 let dropped = Arc::new(AtomicBool::new(false));
4735
4736 let mut player = Player::new();
4737 player.transport = Transport::Loaded(test_session(
4738 QueueItemId::new(),
4739 Box::new(MockEngine {
4740 dropped: dropped.clone(),
4741 }),
4742 ));
4743
4744 player.stop_engine();
4745
4746 assert!(
4750 dropped.load(Ordering::SeqCst),
4751 "AudioEngine must be dropped synchronously in stop_engine (GitHub #89)"
4752 );
4753 }
4754
4755 #[test]
4756 fn abandoning_a_stream_wakes_a_reader_parked_at_the_write_head() {
4757 let live = LiveStream {
4758 feed: crate::remote::downloads::ByteFeed::new(),
4759 abandoned: Default::default(),
4760 };
4761 let feed = live.feed.clone();
4762 let started = std::time::Instant::now();
4763 let reader = thread::spawn(move || {
4764 feed.wait_past(
4765 0,
4766 std::time::Instant::now() + std::time::Duration::from_secs(30),
4767 )
4768 });
4769 thread::sleep(std::time::Duration::from_millis(50));
4770 live.abandon();
4771 reader.join().unwrap();
4772
4773 assert!(live.abandoned.load(std::sync::atomic::Ordering::Acquire));
4774 assert!(started.elapsed() < std::time::Duration::from_secs(5));
4775 }
4776
4777 pub(super) struct StuckBackend {
4782 pub(super) rate: f64,
4783 pub(super) asked: Arc<std::sync::Mutex<Option<(f64, u32)>>>,
4784 pub(super) starts: Arc<std::sync::atomic::AtomicUsize>,
4786 }
4787
4788 struct NullEngine {
4789 starts: Arc<std::sync::atomic::AtomicUsize>,
4790 running: std::sync::atomic::AtomicBool,
4791 lead_in: Arc<AtomicU64>,
4792 }
4793 impl AudioEngineHandle for NullEngine {
4794 fn start(&self) -> Result<(), BackendError> {
4795 self.starts.fetch_add(1, Ordering::Relaxed);
4796 self.running.store(true, Ordering::Relaxed);
4797 Ok(())
4798 }
4799 fn stop(&self) -> Result<(), BackendError> {
4800 self.running.store(false, Ordering::Relaxed);
4801 Ok(())
4802 }
4803 fn is_running(&self) -> bool {
4804 self.running.load(Ordering::Relaxed)
4805 }
4806 fn fade_out(&self) {}
4807 fn fade_in(&self) -> Result<(), BackendError> {
4808 Ok(())
4809 }
4810 fn is_silent(&self) -> bool {
4811 false
4812 }
4813 fn lead_in(&self, frames: u64) {
4814 self.lead_in.store(frames, Ordering::Relaxed);
4815 }
4816 }
4817
4818 impl AudioBackend for StuckBackend {
4819 fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
4820 Ok(vec![self.default_device()?])
4821 }
4822 fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
4823 Ok(backend::DeviceInfo {
4824 name: "Stuck DAC".into(),
4825 sample_rates: vec![self.rate],
4826 platform_id: 0,
4827 kind: Default::default(),
4828 })
4829 }
4830 fn supported_sample_rates(
4831 &self,
4832 _device: &backend::DeviceInfo,
4833 ) -> Result<Vec<f64>, BackendError> {
4834 Ok(vec![self.rate])
4835 }
4836 fn get_device_sample_rate(
4837 &self,
4838 _device: &backend::DeviceInfo,
4839 ) -> Result<f64, BackendError> {
4840 Ok(self.rate)
4841 }
4842 fn set_device_sample_rate(
4843 &self,
4844 _device: &backend::DeviceInfo,
4845 rate: f64,
4846 ) -> Result<f64, BackendError> {
4847 Err(BackendError::UnsupportedSampleRate(rate))
4848 }
4849 fn create_engine(
4850 &self,
4851 _device: &backend::DeviceInfo,
4852 sample_rate: f64,
4853 channels: u32,
4854 _consumer: rtrb::Consumer<f32>,
4855 _samples_played: Arc<AtomicU64>,
4856 ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
4857 *self.asked.lock().unwrap() = Some((sample_rate, channels));
4858 Ok(Box::new(NullEngine {
4859 starts: self.starts.clone(),
4860 running: Default::default(),
4861 lead_in: Default::default(),
4862 }))
4863 }
4864 }
4865
4866 struct CaptureBackend {
4869 consumers: Arc<std::sync::Mutex<Vec<rtrb::Consumer<f32>>>>,
4870 }
4871
4872 impl AudioBackend for CaptureBackend {
4873 fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
4874 Ok(vec![self.default_device()?])
4875 }
4876 fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
4877 Ok(backend::DeviceInfo {
4878 name: "Capture DAC".into(),
4879 sample_rates: vec![44100.0],
4880 platform_id: 0,
4881 kind: Default::default(),
4882 })
4883 }
4884 fn supported_sample_rates(
4885 &self,
4886 _device: &backend::DeviceInfo,
4887 ) -> Result<Vec<f64>, BackendError> {
4888 Ok(vec![44100.0])
4889 }
4890 fn get_device_sample_rate(
4891 &self,
4892 _device: &backend::DeviceInfo,
4893 ) -> Result<f64, BackendError> {
4894 Ok(44100.0)
4895 }
4896 fn set_device_sample_rate(
4897 &self,
4898 _device: &backend::DeviceInfo,
4899 rate: f64,
4900 ) -> Result<f64, BackendError> {
4901 Ok(rate)
4902 }
4903 fn create_engine(
4904 &self,
4905 _device: &backend::DeviceInfo,
4906 _sample_rate: f64,
4907 _channels: u32,
4908 consumer: rtrb::Consumer<f32>,
4909 _samples_played: Arc<AtomicU64>,
4910 ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
4911 self.consumers.lock().unwrap().push(consumer);
4912 Ok(Box::new(NullEngine {
4913 starts: Default::default(),
4914 running: Default::default(),
4915 lead_in: Default::default(),
4916 }))
4917 }
4918 }
4919
4920 fn peak_reaching_the_device(
4922 consumers: &std::sync::Mutex<Vec<rtrb::Consumer<f32>>>,
4923 want: usize,
4924 ) -> f32 {
4925 let mut consumer = consumers.lock().unwrap().pop().expect("an engine was made");
4926 let mut peak = 0.0f32;
4927 let mut got = 0;
4928 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10);
4929 while got < want && std::time::Instant::now() < deadline {
4930 let n = consumer.slots();
4931 if n == 0 {
4932 std::thread::sleep(std::time::Duration::from_millis(2));
4933 continue;
4934 }
4935 let chunk = consumer.read_chunk(n).unwrap();
4936 let (a, b) = chunk.as_slices();
4937 peak = a.iter().chain(b).fold(peak, |p, s| p.max(s.abs()));
4938 got += n;
4939 chunk.commit_all();
4940 }
4941 assert!(got >= want, "only {got} of {want} samples arrived");
4942 peak
4943 }
4944
4945 #[test]
4949 fn a_profile_processes_files_and_streams_alike() {
4950 let dir = tempfile::tempdir().unwrap();
4951 let tone = dir.path().join("tone.wav");
4952 crate::test_utils::generate_wav_tone(&tone, 44100, 440.0, 0.5);
4953 let want = 22050;
4954
4955 let consumers = Arc::new(std::sync::Mutex::new(Vec::new()));
4956 let mut player = Player::new();
4957 player.backend = Box::new(CaptureBackend {
4958 consumers: consumers.clone(),
4959 });
4960 let mut item = make_item("tone");
4961 item.path = tone.clone();
4962 let id = item.id;
4963 player.shared_state.add_items(vec![item]);
4964 player.shared_state.set_cursor(Some(id));
4965
4966 let stream = || {
4967 let feed = crate::remote::downloads::ByteFeed::new();
4968 feed.set(std::fs::metadata(&tone).unwrap().len());
4969 Source::Stream(StreamSource {
4970 path: tone.clone(),
4971 bytes_written: feed,
4972 total: 0,
4973 mode: streaming::ProbeMode::Full,
4974 })
4975 };
4976 let peaks = |player: &mut Player| {
4977 let mut out = Vec::new();
4978 for source in [Source::File(tone.clone()), stream()] {
4979 player
4980 .try_open_session(id, source, None, 0, Run::Playing)
4981 .unwrap();
4982 player.publish();
4983 out.push((
4984 peak_reaching_the_device(&consumers, want),
4985 player.shared_state.dsp().map(|d| d.profile),
4986 ));
4987 player.stop_engine();
4988 }
4989 out
4990 };
4991
4992 let untouched = peaks(&mut player);
4993
4994 player.dsp_override = Some(Arc::new(
4995 crate::audio::dsp::Setup::new(vec![], vec![]).with_preamp(-6.0206),
4996 ));
4997 let processed = peaks(&mut player);
4998
4999 for (kind, ((before, none), (after, half))) in ["file", "stream"]
5000 .iter()
5001 .zip(untouched.iter().zip(&processed))
5002 {
5003 assert_eq!(
5004 none, &None,
5005 "{kind}: nothing is published without a profile"
5006 );
5007 assert_eq!(half.as_deref(), Some("test"), "{kind}: the badge names it");
5008 assert!(*before > 0.1, "{kind}: the tone reached the device");
5009 assert!(
5010 (after / before - 0.5).abs() < 0.01,
5011 "{kind}: {before} → {after}, not halved"
5012 );
5013 }
5014 }
5015
5016 fn engine_format_for(source_rate: u32, channels: u16, device_rate: f64) -> (f64, u32) {
5017 let asked = Arc::new(std::sync::Mutex::new(None));
5018 let mut player = Player::new();
5019 player.backend = Box::new(StuckBackend {
5020 rate: device_rate,
5021 asked: asked.clone(),
5022 starts: Default::default(),
5023 });
5024
5025 let info = buffer::StreamInfo {
5026 codec: "MP3".into(),
5027 sample_rate: source_rate,
5028 channels,
5029 bit_depth: Some(16),
5030 bitrate_kbps: None,
5031 duration_ms: 1000,
5032 };
5033 let (_producer, consumer) = rtrb::RingBuffer::new(16);
5034 player
5035 .create_engine_for(&info, consumer)
5036 .expect("engine creation should succeed");
5037 let asked = *asked.lock().unwrap();
5038 asked.expect("engine was never created")
5039 }
5040
5041 #[test]
5042 fn engine_uses_source_rate_when_device_refuses_switch() {
5043 assert_eq!(engine_format_for(22050, 2, 48000.0), (22050.0, 2));
5046 assert_eq!(engine_format_for(32000, 2, 44100.0), (32000.0, 2));
5047 }
5048
5049 #[test]
5050 fn engine_uses_source_channel_count() {
5051 assert_eq!(engine_format_for(44100, 1, 44100.0), (44100.0, 1));
5052 }
5053
5054 fn settled_rate_for(source_rate: u32, device_rate: f64) -> Option<u32> {
5056 let mut player = Player::new();
5057 player.backend = Box::new(StuckBackend {
5058 rate: device_rate,
5059 asked: Arc::new(std::sync::Mutex::new(None)),
5060 starts: Default::default(),
5061 });
5062 let state = player.shared_state.clone();
5063
5064 let info = buffer::StreamInfo {
5065 codec: "MP3".into(),
5066 sample_rate: source_rate,
5067 channels: 2,
5068 bit_depth: Some(16),
5069 bitrate_kbps: None,
5070 duration_ms: 1000,
5071 };
5072 let (_producer, consumer) = rtrb::RingBuffer::new(16);
5073 player
5074 .create_engine_for(&info, consumer)
5075 .expect("engine creation should succeed");
5076 state.output_sample_rate()
5077 }
5078
5079 #[test]
5080 fn settled_device_rate_reaches_the_shared_state() {
5081 assert_eq!(settled_rate_for(22050, 48000.0), Some(48000));
5085 assert_eq!(settled_rate_for(44100, 44100.0), Some(44100));
5087 }
5088
5089 struct SlowBackend {
5091 observed: Arc<std::sync::Mutex<Vec<Option<u32>>>>,
5092 state: Arc<SharedPlayerState>,
5093 lead_in: Arc<AtomicU64>,
5094 }
5095
5096 impl AudioBackend for SlowBackend {
5097 fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
5098 Ok(vec![self.default_device()?])
5099 }
5100 fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
5101 Ok(backend::DeviceInfo {
5102 name: "Slow DAC".into(),
5103 sample_rates: vec![44100.0, 48000.0],
5104 platform_id: 0,
5105 kind: Default::default(),
5106 })
5107 }
5108 fn supported_sample_rates(
5109 &self,
5110 _device: &backend::DeviceInfo,
5111 ) -> Result<Vec<f64>, BackendError> {
5112 Ok(vec![44100.0, 48000.0])
5113 }
5114 fn get_device_sample_rate(
5115 &self,
5116 _device: &backend::DeviceInfo,
5117 ) -> Result<f64, BackendError> {
5118 Ok(48000.0)
5119 }
5120 fn set_device_sample_rate(
5121 &self,
5122 _device: &backend::DeviceInfo,
5123 rate: f64,
5124 ) -> Result<f64, BackendError> {
5125 self.observed
5127 .lock()
5128 .unwrap()
5129 .push(self.state.output_sample_rate());
5130 Ok(rate)
5131 }
5132 fn create_engine(
5133 &self,
5134 _device: &backend::DeviceInfo,
5135 _sample_rate: f64,
5136 _channels: u32,
5137 _consumer: rtrb::Consumer<f32>,
5138 _samples_played: Arc<AtomicU64>,
5139 ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
5140 Ok(Box::new(NullEngine {
5141 starts: Default::default(),
5142 running: Default::default(),
5143 lead_in: self.lead_in.clone(),
5144 }))
5145 }
5146 }
5147
5148 #[test]
5149 fn the_previous_rate_is_not_published_while_the_device_reclocks() {
5150 let mut player = Player::new();
5156 let state = player.shared_state.clone();
5157 state.set_output_sample_rate(48000);
5158
5159 let observed = Arc::new(std::sync::Mutex::new(Vec::new()));
5160 player.backend = Box::new(SlowBackend {
5161 observed: observed.clone(),
5162 state: state.clone(),
5163 lead_in: Default::default(),
5164 });
5165
5166 let info = buffer::StreamInfo {
5167 codec: "FLAC".into(),
5168 sample_rate: 44100,
5169 channels: 2,
5170 bit_depth: Some(16),
5171 bitrate_kbps: None,
5172 duration_ms: 1000,
5173 };
5174 let (_producer, consumer) = rtrb::RingBuffer::new(16);
5175 player
5176 .create_engine_for(&info, consumer)
5177 .expect("engine creation should succeed");
5178
5179 assert_eq!(
5180 *observed.lock().unwrap(),
5181 vec![None],
5182 "mid-switch the output rate must read as unknown, not as the last track's"
5183 );
5184 assert_eq!(state.output_sample_rate(), Some(44100));
5185 }
5186
5187 fn lead_in_for(source_rate: u32) -> u64 {
5190 let mut player = Player::new();
5191 let lead_in = Arc::new(AtomicU64::new(0));
5192 player.backend = Box::new(SlowBackend {
5193 observed: Default::default(),
5194 state: player.shared_state.clone(),
5195 lead_in: lead_in.clone(),
5196 });
5197 let info = buffer::StreamInfo {
5198 codec: "FLAC".into(),
5199 sample_rate: source_rate,
5200 channels: 2,
5201 bit_depth: Some(16),
5202 bitrate_kbps: None,
5203 duration_ms: 1000,
5204 };
5205 let (_producer, consumer) = rtrb::RingBuffer::new(16);
5206 player
5207 .create_engine_for(&info, consumer)
5208 .expect("engine creation should succeed");
5209 lead_in.load(Ordering::Relaxed)
5210 }
5211
5212 #[test]
5213 fn an_engine_made_inside_the_silence_keeps_the_rest_of_it() {
5214 let mut player = Player::new();
5217 let lead_in = Arc::new(AtomicU64::new(0));
5218 player.backend = Box::new(SlowBackend {
5219 observed: Default::default(),
5220 state: player.shared_state.clone(),
5221 lead_in: lead_in.clone(),
5222 });
5223 let info = |sample_rate| buffer::StreamInfo {
5224 codec: "FLAC".into(),
5225 sample_rate,
5226 channels: 2,
5227 bit_depth: Some(16),
5228 bitrate_kbps: None,
5229 duration_ms: 1000,
5230 };
5231 let (_p, consumer) = rtrb::RingBuffer::new(16);
5232 player.create_engine_for(&info(44100), consumer).unwrap();
5233 let first = lead_in.swap(0, Ordering::Relaxed);
5234 assert!(first > 0);
5235
5236 let (_p, consumer) = rtrb::RingBuffer::new(16);
5237 player.create_engine_for(&info(48000), consumer).unwrap();
5238 let carried = lead_in.load(Ordering::Relaxed);
5239 assert!(
5241 carried > 0 && carried < first * 48000 / 44100,
5242 "carried {carried} of {first}"
5243 );
5244 }
5245
5246 #[test]
5247 fn silence_leads_in_only_after_the_device_changed_rate() {
5248 assert!(
5249 lead_in_for(44100) > 0,
5250 "the device is relocking, so the start of the track would be lost"
5251 );
5252 assert_eq!(lead_in_for(48000), 0, "no switch, nothing to wait for");
5253 }
5254
5255 struct WatchedBackend {
5257 inner: StuckBackend,
5258 #[allow(clippy::type_complexity)]
5259 captured: Arc<std::sync::Mutex<Option<Box<dyn Fn(f64) + Send + Sync>>>>,
5260 }
5261
5262 struct NullWatch;
5263 impl backend::SampleRateWatch for NullWatch {}
5264
5265 impl AudioBackend for WatchedBackend {
5266 fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
5267 self.inner.list_devices()
5268 }
5269 fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
5270 self.inner.default_device()
5271 }
5272 fn supported_sample_rates(
5273 &self,
5274 device: &backend::DeviceInfo,
5275 ) -> Result<Vec<f64>, BackendError> {
5276 self.inner.supported_sample_rates(device)
5277 }
5278 fn get_device_sample_rate(
5279 &self,
5280 device: &backend::DeviceInfo,
5281 ) -> Result<f64, BackendError> {
5282 self.inner.get_device_sample_rate(device)
5283 }
5284 fn set_device_sample_rate(
5285 &self,
5286 device: &backend::DeviceInfo,
5287 rate: f64,
5288 ) -> Result<f64, BackendError> {
5289 self.inner.set_device_sample_rate(device, rate)
5290 }
5291 fn watch_device_sample_rate(
5292 &self,
5293 _device: &backend::DeviceInfo,
5294 on_change: Box<dyn Fn(f64) + Send + Sync>,
5295 ) -> Option<Box<dyn backend::SampleRateWatch>> {
5296 *self.captured.lock().unwrap() = Some(on_change);
5297 Some(Box::new(NullWatch))
5298 }
5299 fn create_engine(
5300 &self,
5301 device: &backend::DeviceInfo,
5302 sample_rate: f64,
5303 channels: u32,
5304 consumer: rtrb::Consumer<f32>,
5305 samples_played: Arc<AtomicU64>,
5306 ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
5307 self.inner
5308 .create_engine(device, sample_rate, channels, consumer, samples_played)
5309 }
5310 }
5311
5312 #[test]
5313 fn external_rate_change_reaches_the_shared_state() {
5314 let captured = Arc::new(std::sync::Mutex::new(None));
5318 let mut player = Player::new();
5319 player.backend = Box::new(WatchedBackend {
5320 inner: StuckBackend {
5321 rate: 44100.0,
5322 asked: Arc::new(std::sync::Mutex::new(None)),
5323 starts: Default::default(),
5324 },
5325 captured: captured.clone(),
5326 });
5327 let state = player.shared_state.clone();
5328
5329 let info = buffer::StreamInfo {
5330 codec: "FLAC".into(),
5331 sample_rate: 44100,
5332 channels: 2,
5333 bit_depth: Some(16),
5334 bitrate_kbps: None,
5335 duration_ms: 1000,
5336 };
5337 let (_producer, consumer) = rtrb::RingBuffer::new(16);
5338 player
5339 .create_engine_for(&info, consumer)
5340 .expect("engine creation should succeed");
5341 assert_eq!(state.output_sample_rate(), Some(44100));
5342
5343 let on_change = captured.lock().unwrap().take().expect("watch registered");
5344 on_change(48000.0);
5345 assert_eq!(state.output_sample_rate(), Some(48000));
5346 }
5347}