1pub mod commands;
2pub mod history;
3pub mod state;
4pub mod undo;
5
6use std::path::{Path, PathBuf};
7use std::sync::Arc;
8use std::thread;
9
10use thiserror::Error;
11
12use crate::audio::{
13 analyzer::VizAnalyzer,
14 backend::{self, AudioBackend, AudioEngineHandle, BackendError, SampleRateWatch},
15 buffer, streaming,
16 viz::{VizBuffer, VizSnapshot},
17};
18use crate::remote::client::PlaybackReportState;
19use buffer::PlaybackTimeline;
20use commands::{CommandChannel, PlayerCommand};
21use history::{InFlight, PlayEvent, PlayRecorder, PlaybackReport};
22use state::{
23 ItemState, LoadState, PlaybackSource, PlaybackState, QueueItemId, SharedPlayerState, TrackInfo,
24};
25use undo::{UndoEntry, UndoStack};
26
27pub(crate) const RING_BUFFER_SIZE: usize = 192_000 * 2;
29
30const SEEK_END_GUARD_MS: u64 = 500;
33
34const BOUNDARY_SLACK: std::time::Duration = std::time::Duration::from_millis(5);
37
38const FADE_CHECK: std::time::Duration = std::time::Duration::from_millis(50);
41
42#[derive(Debug, Error)]
43pub enum PlayerError {
44 #[error("backend error: {0}")]
45 Backend(#[from] BackendError),
46 #[error("decode error: {0}")]
47 Decode(#[from] buffer::DecodeError),
48}
49
50#[derive(Clone)]
53struct StreamSource {
54 path: PathBuf,
55 bytes_written: Arc<crate::remote::downloads::ByteFeed>,
56 total: u64,
57 mode: streaming::ProbeMode,
58}
59
60fn media_extension(path: &Path) -> Option<String> {
63 crate::remote::download::strip_part_suffix(path)
64 .extension()
65 .and_then(|e| e.to_str())
66 .map(str::to_ascii_lowercase)
67}
68
69fn hint_for(path: &Path) -> symphonia::core::formats::probe::Hint {
71 let mut hint = symphonia::core::formats::probe::Hint::new();
72 if let Some(ext) = media_extension(path) {
73 hint.with_extension(&ext);
74 }
75 hint
76}
77
78fn lengthless_mode_for(path: &Path) -> streaming::ProbeMode {
88 match media_extension(path).as_deref() {
89 Some("ogg" | "oga" | "opus" | "spx" | "m4a" | "m4b" | "mp4" | "mov") => {
90 streaming::ProbeMode::LengthlessWholeEnd
91 }
92 _ => streaming::ProbeMode::Lengthless,
93 }
94}
95
96pub struct Player {
98 shared_state: Arc<SharedPlayerState>,
99 commands: CommandChannel,
100 active_playback: Option<ActivePlayback>,
101 timeline: Arc<PlaybackTimeline>,
102 viz_buffer: Arc<VizBuffer>,
103 viz_snapshot: Arc<VizSnapshot>,
104 _viz_analyzer: VizAnalyzer,
106 undo_stack: UndoStack,
107 batch_buffer: Option<Vec<UndoEntry>>,
110 output_device_name: Option<String>,
112 backend: Box<dyn AudioBackend>,
114 last_skip: std::time::Instant,
116 stream_mode: streaming::ProbeMode,
119 history: Option<PlayRecorder>,
122 in_flight: Option<InFlight>,
124 lead_in_ends: Option<std::time::Instant>,
126 #[cfg(test)]
129 playback_starts: usize,
130}
131
132#[derive(Clone, Copy, PartialEq, Eq)]
134enum Start {
135 Playing,
136 Paused,
137}
138
139struct ActivePlayback {
141 engine: Box<dyn AudioEngineHandle>,
142 decode_handle: buffer::DecodeHandle,
143 stream: Option<LiveStream>,
145 _rate_watch: Option<Box<dyn SampleRateWatch>>,
148}
149
150struct LiveStream {
154 feed: Arc<crate::remote::downloads::ByteFeed>,
155 abandoned: Arc<std::sync::atomic::AtomicBool>,
156}
157
158impl LiveStream {
159 fn abandon(&self) {
160 self.abandoned
161 .store(true, std::sync::atomic::Ordering::Release);
162 self.feed.done();
163 }
164}
165
166impl Default for Player {
167 fn default() -> Self {
168 Self::new()
169 }
170}
171
172impl Player {
173 pub fn new() -> Self {
174 let viz_buffer = VizBuffer::new();
175 let viz_snapshot = VizSnapshot::new();
176 let timeline = PlaybackTimeline::new();
177 let cfg = crate::config::Config::load_or_default();
178 let viz_analyzer = VizAnalyzer::spawn_with_snapshot(
179 Arc::clone(&viz_buffer),
180 &cfg.visualizer,
181 Arc::clone(&viz_snapshot),
182 timeline.samples_played_counter(),
183 );
184
185 let shared_state = SharedPlayerState::new();
186 shared_state.attach_timeline(timeline.clone());
187 let commands = CommandChannel::new();
188 let tx = commands.tx.clone();
189 timeline.on_queued(move || {
190 let _ = tx.try_send(PlayerCommand::TrackQueued);
191 });
192
193 Self {
194 shared_state,
195 commands,
196 active_playback: None,
197 lead_in_ends: None,
198 timeline,
199 viz_buffer,
200 viz_snapshot,
201 _viz_analyzer: viz_analyzer,
202 undo_stack: UndoStack::new(),
203 batch_buffer: None,
204 output_device_name: cfg.playback.output_device,
205 backend: crate::audio::platform_backend(),
206 last_skip: std::time::Instant::now(),
207 stream_mode: streaming::ProbeMode::Full,
208 history: None,
209 in_flight: None,
210 #[cfg(test)]
211 playback_starts: 0,
212 }
213 }
214
215 pub fn shared_state(&self) -> Arc<SharedPlayerState> {
217 self.shared_state.clone()
218 }
219
220 pub fn timeline(&self) -> Arc<PlaybackTimeline> {
222 self.timeline.clone()
223 }
224
225 pub fn viz_buffer(&self) -> Arc<VizBuffer> {
227 self.viz_buffer.clone()
228 }
229
230 pub fn viz_snapshot(&self) -> Arc<VizSnapshot> {
233 self.viz_snapshot.clone()
234 }
235
236 pub fn undo_stack(&self) -> &UndoStack {
238 &self.undo_stack
239 }
240
241 #[allow(clippy::type_complexity)]
249 fn create_engine_for(
250 &mut self,
251 info: &buffer::StreamInfo,
252 consumer: rtrb::Consumer<f32>,
253 ) -> Result<(Box<dyn AudioEngineHandle>, Option<Box<dyn SampleRateWatch>>), PlayerError> {
254 let device = self.resolve_device()?;
255 let device_rate = self.backend.get_device_sample_rate(&device)?;
256 let source_rate = info.sample_rate as f64;
257
258 self.shared_state.clear_output_sample_rate();
262
263 let settled = if (device_rate - source_rate).abs() > 0.1 {
264 log::info!(
265 "switching device sample rate: {}Hz → {}Hz",
266 device_rate,
267 source_rate
268 );
269 match self.backend.set_device_sample_rate(&device, source_rate) {
270 Ok(rate) => rate,
271 Err(e) => {
272 log::warn!("failed to set device sample rate: {}", e);
273 device_rate
274 }
275 }
276 } else {
277 device_rate
278 };
279
280 if (settled - source_rate).abs() > 0.1 {
281 log::warn!(
282 "device stayed at {}Hz (wanted {}Hz) — output is resampled, not bit-perfect",
283 settled,
284 source_rate
285 );
286 }
287
288 self.shared_state
291 .set_output_sample_rate(settled.round() as u32);
292
293 let watch_state = self.shared_state.clone();
297 let watch_name = device.name.clone();
298 let rate_watch = self.backend.watch_device_sample_rate(
299 &device,
300 Box::new(move |rate| {
301 log::info!("device sample rate changed externally: {rate}Hz on '{watch_name}'");
302 watch_state.set_output_sample_rate(rate.round() as u32);
303 }),
304 );
305
306 let engine = self.backend.create_engine(
307 &device,
308 source_rate,
309 info.channels as u32,
310 consumer,
311 self.timeline.samples_played_counter(),
312 )?;
313 let now = std::time::Instant::now();
317 let lead_in = if device_rate > 0.0 && (settled - device_rate).abs() > 0.1 {
318 std::time::Duration::from_millis(
319 crate::config::Config::load_or_default()
320 .playback
321 .rate_switch_lead_in_ms as u64,
322 )
323 } else {
324 self.lead_in_ends
325 .map(|end| end.saturating_duration_since(now))
326 .unwrap_or_default()
327 };
328 self.lead_in_ends = None;
329 if !lead_in.is_zero() {
330 engine.lead_in((settled * lead_in.as_secs_f64()) as u64);
331 self.lead_in_ends = Some(now + lead_in);
332 }
333
334 Ok((engine, rate_watch))
335 }
336
337 fn resolve_device(&self) -> Result<backend::DeviceInfo, PlayerError> {
340 if let Some(ref name) = self.output_device_name {
341 match self.backend.list_devices() {
342 Ok(devices) => {
343 if let Some(dev) = devices.into_iter().find(|d| d.name == *name) {
344 return Ok(dev);
345 }
346 log::warn!(
347 "configured output device '{}' not found, falling back to default",
348 name,
349 );
350 }
351 Err(e) => {
352 log::warn!("failed to list devices while resolving '{}': {}", name, e);
353 }
354 }
355 }
356 Ok(self.backend.default_device()?)
357 }
358
359 pub fn set_output_device(&mut self, name: String) {
362 log::info!("switching output device to: {}", name);
363 self.output_device_name = Some(name.clone());
364
365 if let Err(e) = crate::config::Config::persist(|cfg| {
366 cfg.playback.output_device = Some(name);
367 }) {
368 log::error!("failed to save output device config: {}", e);
369 }
370
371 self.restart_on_current_track();
372 }
373
374 pub fn clear_output_device(&mut self) {
376 log::info!("reverting to system default output device");
377 self.output_device_name = None;
378
379 if let Err(e) = crate::config::Config::persist(|cfg| {
380 cfg.playback.output_device = None;
381 }) {
382 log::error!("failed to save output device config: {}", e);
383 }
384
385 self.restart_on_current_track();
386 }
387
388 fn restart_on_current_track(&mut self) {
391 if let Some(info) = self.shared_state.track_info() {
392 let position_ms = self.shared_state.position_ms();
393 if let Err(e) = self.restart_current(&info, position_ms) {
394 log::error!("failed to restart playback on device switch: {}", e);
395 }
396 }
397 }
398
399 pub fn output_device_name(&self) -> Option<&str> {
401 self.output_device_name.as_deref()
402 }
403
404 pub fn command_sender(&self) -> crossbeam_channel::Sender<PlayerCommand> {
406 self.commands.tx.clone()
407 }
408
409 pub fn play(&mut self, id: QueueItemId) {
412 self.shared_state.set_cursor(Some(id));
413
414 match self.shared_state.item_playback_source(id) {
415 Some(PlaybackSource::Ready(path)) => {
416 if let Err(e) = self.start_playback(id, &path, 0, Start::Playing) {
417 log::error!("play failed: {}", e);
418 }
419 }
420 Some(PlaybackSource::Streaming {
421 path,
422 bytes_written,
423 total,
424 }) => {
425 self.report(PlaybackReportState::Stopped);
429 self.stop_engine();
430 self.shared_state.set_playback_state(PlaybackState::Stopped);
431 self.probe_stream_for_playback(id, &path, bytes_written, total);
432 }
433 None => {
434 self.report(PlaybackReportState::Stopped);
436 self.stop_engine();
437 self.shared_state.set_playback_state(PlaybackState::Stopped);
438 log::info!("play: item {:?} not ready, waiting for TrackReady", id);
439 }
440 }
441 }
442
443 fn cue(&mut self, id: QueueItemId, position_ms: u64, start: Start) {
447 self.shared_state.set_cursor(Some(id));
448 let Some(PlaybackSource::Ready(path)) = self.shared_state.item_playback_source(id) else {
449 if start == Start::Playing {
450 self.play(id);
451 }
452 return;
453 };
454 if let Err(e) = self.start_playback(id, &path, position_ms, start) {
455 log::error!("cue failed: {}", e);
456 return;
457 }
458 self.report(match start {
459 Start::Playing => PlaybackReportState::Playing,
460 Start::Paused => PlaybackReportState::Paused,
461 });
462 }
463
464 fn start_playback(
469 &mut self,
470 id: QueueItemId,
471 path: &Path,
472 seek_ms: u64,
473 start: Start,
474 ) -> Result<(), PlayerError> {
475 #[cfg(test)]
476 {
477 self.playback_starts += 1;
478 }
479 let result = self.open_playback(id, path, seek_ms, start);
480 if result.is_err() {
481 self.stop_playback_and_clear_state();
482 }
483 self.wake_analyzer();
484 result
485 }
486
487 fn open_playback(
488 &mut self,
489 id: QueueItemId,
490 path: &Path,
491 seek_ms: u64,
492 start: Start,
493 ) -> Result<(), PlayerError> {
494 self.stop_engine();
495
496 let info = buffer::probe_file(path)?;
497
498 self.shared_state.set_track_info(Some(TrackInfo {
502 id,
503 path: path.to_path_buf(),
504 codec: info.codec.clone(),
505 sample_rate: info.sample_rate,
506 bit_depth: info.bit_depth,
507 bitrate_kbps: info.bitrate_kbps,
508 channels: info.channels,
509 duration_ms: info.duration_ms,
510 }));
511 self.shared_state.set_position_ms(seek_ms);
512 self.on_track_changed(id, seek_ms);
513 log::info!(
514 "playing: {} ({:?}) — {} {}Hz/{}ch, {}ms{}",
515 path.display(),
516 id,
517 info.codec,
518 info.sample_rate,
519 info.channels,
520 info.duration_ms,
521 if seek_ms > 0 {
522 format!(" @{}ms", seek_ms)
523 } else {
524 String::new()
525 }
526 );
527
528 let (producer, consumer) = rtrb::RingBuffer::new(RING_BUFFER_SIZE);
529
530 self.timeline.reset();
531
532 let next_track = self.decode_cursor(id);
533
534 let cfg = crate::config::Config::load_or_default();
536 let rg_mode = cfg.playback.replaygain;
537 let pre_amp_db = cfg.playback.pre_amp_db;
538
539 let finish_tx = self.commands.tx.clone();
540 let (_stream_info, decode_handle) = buffer::start_decode_file(
541 id,
542 path,
543 producer,
544 seek_ms,
545 next_track,
546 self.timeline.clone(),
547 Some(self.viz_buffer.clone()),
548 rg_mode,
549 pre_amp_db,
550 move || {
551 finish_tx.send(PlayerCommand::DecodeFinished).ok();
552 },
553 )?;
554
555 let (engine, rate_watch) = self.create_engine_for(&info, consumer)?;
556 let state = match start {
557 Start::Playing => {
558 engine.start()?;
559 PlaybackState::Playing
560 }
561 Start::Paused => PlaybackState::Paused,
564 };
565 self.shared_state.set_playback_state(state);
566
567 self.active_playback = Some(ActivePlayback {
568 engine,
569 decode_handle,
570 stream: None,
571 _rate_watch: rate_watch,
572 });
573
574 Ok(())
575 }
576
577 fn decode_cursor(
581 &self,
582 id: QueueItemId,
583 ) -> impl Fn() -> Option<(QueueItemId, PathBuf)> + Send + 'static {
584 let state = self.shared_state.clone();
585 let cursor = parking_lot::Mutex::new(Some(id));
586 move || {
587 let current = cursor.lock().take()?;
588 let next = state.peek_next_ready_after(current);
589 if let Some((next_id, _)) = &next {
590 *cursor.lock() = Some(*next_id);
591 }
592 next
593 }
594 }
595
596 fn probe_stream_for_playback(
611 &self,
612 id: QueueItemId,
613 path: &Path,
614 bytes_written: Arc<crate::remote::downloads::ByteFeed>,
615 total: u64,
616 ) {
617 let path = path.to_path_buf();
618 let tx = self.commands.tx.clone();
619 let hint = hint_for(&path);
620
621 let status = {
626 let downloading = self.stream_status_fn(id);
627 let state = self.shared_state.clone();
628 Arc::new(move || {
629 if state.is_cursor(id) {
630 downloading()
631 } else {
632 streaming::StreamStatus::Failed
633 }
634 }) as Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync>
635 };
636
637 let spawned = thread::Builder::new()
638 .name("koan-stream-probe".into())
639 .spawn(move || {
640 let attempt = |mode, wait: bool| {
647 let open = if wait {
648 streaming::PartialFileSource::open(
649 &path,
650 bytes_written.clone(),
651 total,
652 status.clone(),
653 mode,
654 )
655 } else {
656 streaming::PartialFileSource::open_for_probe(
657 &path,
658 bytes_written.clone(),
659 total,
660 status.clone(),
661 mode,
662 )
663 };
664 open.map_err(buffer::DecodeError::Io).and_then(|source| {
665 let mss = symphonia::core::io::MediaSourceStream::new(
666 Box::new(source),
667 Default::default(),
668 );
669 buffer::probe_source(mss, &hint)
670 })
671 };
672
673 let info = match attempt(streaming::ProbeMode::Full, false) {
677 Ok(info) => Some((info, streaming::ProbeMode::Full)),
678 Err(e) => {
679 log::info!(
685 "stream probe: {} needs more than has arrived ({}), opening without a length",
686 path.display(),
687 e
688 );
689 let lengthless = lengthless_mode_for(&path);
690 attempt(lengthless, true)
691 .ok()
692 .map(|info| (info, lengthless))
693 }
694 };
695
696 match info {
697 Some((info, mode)) => {
698 tx.send(PlayerCommand::StreamProbed {
699 id,
700 info: Box::new(info),
701 mode,
702 })
703 .ok();
704 }
705 None => log::info!(
708 "stream probe: {} cannot start early, waiting for the download",
709 path.display()
710 ),
711 }
712 });
713
714 if let Err(e) = spawned {
715 log::warn!("stream probe: could not spawn for {:?}: {}", id, e);
716 }
717 }
718
719 fn stream_probed(
722 &mut self,
723 id: QueueItemId,
724 info: buffer::StreamInfo,
725 mode: streaming::ProbeMode,
726 ) {
727 if !self.shared_state.is_cursor(id) {
728 return; }
730 if self.shared_state.playback_state() != PlaybackState::Stopped {
731 return; }
733
734 match self.shared_state.item_playback_source(id) {
735 Some(PlaybackSource::Ready(path)) => {
737 if let Err(e) = self.start_playback(id, &path, 0, Start::Playing) {
738 log::error!("stream probe: playback failed: {}", e);
739 }
740 }
741 Some(PlaybackSource::Streaming {
742 path,
743 bytes_written,
744 total,
745 }) => {
746 let source = StreamSource {
747 path,
748 bytes_written,
749 total,
750 mode,
751 };
752 if let Err(e) = self.start_streaming_playback(id, source, 0, info, Start::Playing) {
753 log::error!("stream probe: streaming playback failed: {}", e);
754 }
755 }
756 None => {}
757 }
758 }
759
760 fn stream_status_fn(
764 &self,
765 id: QueueItemId,
766 ) -> Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync> {
767 let state = self.shared_state.clone();
768 Arc::new(move || match state.item_load_state(id) {
769 Some(LoadState::Ready) => streaming::StreamStatus::Complete,
770 Some(LoadState::Failed(_)) => streaming::StreamStatus::Failed,
771 _ => streaming::StreamStatus::Downloading,
772 })
773 }
774
775 fn start_streaming_playback(
785 &mut self,
786 id: QueueItemId,
787 source: StreamSource,
788 seek_ms: u64,
789 info: buffer::StreamInfo,
790 start: Start,
791 ) -> Result<(), PlayerError> {
792 let result = self.open_streaming_playback(id, source, seek_ms, info, start);
793 if result.is_err() {
794 self.stop_playback_and_clear_state();
795 }
796 result
797 }
798
799 fn open_streaming_playback(
800 &mut self,
801 id: QueueItemId,
802 source: StreamSource,
803 seek_ms: u64,
804 info: buffer::StreamInfo,
805 start: Start,
806 ) -> Result<(), PlayerError> {
807 self.stop_engine();
808 self.stream_mode = source.mode;
810 let path = source.path.as_path();
811
812 let live = LiveStream {
813 feed: source.bytes_written.clone(),
814 abandoned: Default::default(),
815 };
816 let status = {
817 let downloading = self.stream_status_fn(id);
818 let abandoned = live.abandoned.clone();
819 Arc::new(move || {
820 if abandoned.load(std::sync::atomic::Ordering::Acquire) {
821 streaming::StreamStatus::Failed
822 } else {
823 downloading()
824 }
825 }) as Arc<dyn Fn() -> streaming::StreamStatus + Send + Sync>
826 };
827 let open_source = {
828 let StreamSource {
829 path,
830 bytes_written,
831 total,
832 mode,
833 } = source.clone();
834 let status = status.clone();
835 move || {
836 streaming::PartialFileSource::open(
837 &path,
838 bytes_written.clone(),
839 total,
840 status.clone(),
841 mode,
842 )
843 }
844 };
845
846 self.shared_state.set_track_info(Some(TrackInfo {
847 id,
848 path: path.to_path_buf(),
849 codec: info.codec.clone(),
850 sample_rate: info.sample_rate,
851 bit_depth: info.bit_depth,
852 bitrate_kbps: info.bitrate_kbps,
853 channels: info.channels,
854 duration_ms: info.duration_ms,
855 }));
856 self.shared_state.set_position_ms(seek_ms);
857 self.on_track_changed(id, seek_ms);
858 log::info!(
859 "streaming: {} ({:?}) — {} {}Hz/{}ch, {}ms{}",
860 path.display(),
861 id,
862 info.codec,
863 info.sample_rate,
864 info.channels,
865 info.duration_ms,
866 if seek_ms > 0 {
867 format!(" @{}ms", seek_ms)
868 } else {
869 String::new()
870 },
871 );
872
873 let (producer, consumer) = rtrb::RingBuffer::new(RING_BUFFER_SIZE);
874
875 self.timeline.reset();
876
877 let next_track = self.decode_cursor(id);
879
880 let first = buffer::SourceEntry {
881 id,
882 path: path.to_path_buf(),
883 hint: hint_for(path),
884 make_mss: Box::new(move || {
885 Ok(symphonia::core::io::MediaSourceStream::new(
886 Box::new(open_source()?),
887 Default::default(),
888 ))
889 }),
890 };
891
892 let cfg = crate::config::Config::load_or_default();
894 let rg_mode = cfg.playback.replaygain;
895 let pre_amp_db = cfg.playback.pre_amp_db;
896
897 let finish_tx = self.commands.tx.clone();
898 let (_stream_info, decode_handle) = buffer::start_decode(
899 first,
900 producer,
901 seek_ms,
902 move || {
903 let (next_id, next_path) = next_track()?;
904 Some(buffer::SourceEntry::from_file(next_id, next_path))
905 },
906 self.timeline.clone(),
907 Some(self.viz_buffer.clone()),
908 rg_mode,
909 pre_amp_db,
910 move || {
911 finish_tx.send(PlayerCommand::DecodeFinished).ok();
912 },
913 )?;
914
915 let (engine, rate_watch) = self.create_engine_for(&info, consumer)?;
916 let state = match start {
917 Start::Playing => {
918 engine.start()?;
919 PlaybackState::Playing
920 }
921 Start::Paused => PlaybackState::Paused,
924 };
925 self.shared_state.set_playback_state(state);
926
927 self.active_playback = Some(ActivePlayback {
928 engine,
929 decode_handle,
930 stream: Some(live),
931 _rate_watch: rate_watch,
932 });
933
934 Ok(())
935 }
936
937 pub fn seek(&mut self, position_ms: u64) {
944 let Some(info) = self.shared_state.track_info() else {
945 return;
946 };
947 let seekable = self.shared_state.seekable_ms();
949 if seekable == 0 {
950 log::debug!("seek declined: {:?} is not seekable yet", info.id);
954 return;
955 }
956 let ceiling = seekable.min(
957 self.shared_state
958 .duration_ms()
959 .saturating_sub(SEEK_END_GUARD_MS),
960 );
961 let clamped = position_ms.min(ceiling);
962
963 if let Err(e) = self.restart_current(&info, clamped) {
964 log::error!("seek failed: {}", e);
965 }
966 }
967
968 fn restart_current(&mut self, info: &TrackInfo, position_ms: u64) -> Result<(), PlayerError> {
976 let was_paused = self.shared_state.playback_state() == PlaybackState::Paused;
977 let start = if was_paused {
978 Start::Paused
979 } else {
980 Start::Playing
981 };
982
983 match self.shared_state.item_playback_source(info.id) {
984 Some(PlaybackSource::Streaming {
985 path,
986 bytes_written,
987 total,
988 }) => {
989 let known = buffer::StreamInfo {
993 codec: info.codec.clone(),
994 sample_rate: info.sample_rate,
995 channels: info.channels,
996 bit_depth: info.bit_depth,
997 bitrate_kbps: info.bitrate_kbps,
998 duration_ms: info.duration_ms,
999 };
1000 let source = StreamSource {
1001 path,
1002 bytes_written,
1003 total,
1004 mode: self.stream_mode,
1005 };
1006 self.start_streaming_playback(info.id, source, position_ms, known, start)?;
1007 }
1008 Some(PlaybackSource::Ready(path)) => {
1009 self.start_playback(info.id, &path, position_ms, start)?;
1010 }
1011 None => return Ok(()),
1012 }
1013
1014 self.report(if was_paused {
1015 PlaybackReportState::Paused
1016 } else {
1017 PlaybackReportState::Playing
1018 });
1019 Ok(())
1020 }
1021
1022 pub fn next_track(&mut self) {
1024 match self.shared_state.advance_cursor_loadable() {
1025 Some(id) => self.play(id),
1026 None => {
1027 log::info!("no more tracks in playlist");
1028 self.stop_playback_and_clear_state();
1029 }
1030 }
1031 }
1032
1033 pub fn prev_track(&mut self) {
1035 match self.shared_state.retreat_cursor() {
1036 Some((id, _)) => self.play(id),
1037 None => {
1038 if let Some(info) = self.shared_state.track_info()
1040 && let Err(e) = self.restart_current(&info, 0)
1041 {
1042 log::error!("restart failed: {}", e);
1043 }
1044 }
1045 }
1046 }
1047
1048 pub fn pause(&mut self) {
1053 let Some(ref playback) = self.active_playback else {
1054 return;
1055 };
1056 if crate::config::Config::load_or_default()
1057 .playback
1058 .fade_on_pause
1059 {
1060 playback.engine.fade_out();
1061 self.shared_state.set_playback_state(PlaybackState::Paused);
1062 } else {
1063 self.pause_now();
1064 }
1065 self.report(PlaybackReportState::Paused);
1066 }
1067
1068 fn pause_now(&mut self) {
1071 if let Some(ref playback) = self.active_playback {
1072 if let Err(e) = playback.engine.stop() {
1073 log::error!("pause failed: {}", e);
1074 return;
1075 }
1076 self.shared_state.set_playback_state(PlaybackState::Paused);
1077 }
1078 }
1079
1080 pub fn resume(&mut self) {
1086 if self.active_playback.is_none() {
1087 if let Some(id) = self.shared_state.cursor() {
1088 self.play(id);
1089 }
1090 return;
1091 }
1092 if let Some(ref playback) = self.active_playback {
1093 let engine = &playback.engine;
1094 let resumed = if engine.is_running() || engine.is_silent() {
1095 self.lead_in_ends = None;
1096 engine.fade_in()
1097 } else {
1098 engine.start()
1099 };
1100 if let Err(e) = resumed {
1101 log::error!("resume failed: {}", e);
1102 return;
1103 }
1104 self.shared_state.set_playback_state(PlaybackState::Playing);
1105 self.wake_analyzer();
1106 self.report(PlaybackReportState::Playing);
1107 }
1108 }
1109
1110 fn wake_analyzer(&self) {
1117 self.viz_snapshot.wake();
1118 }
1119
1120 pub fn stop(&mut self) {
1122 self.shared_state.clear_playlist();
1123 self.stop_playback_and_clear_state();
1124 }
1125
1126 fn stop_engine(&mut self) {
1132 let Some(playback) = self.active_playback.take() else {
1133 return;
1134 };
1135 self.bank_listening();
1136 let ActivePlayback {
1137 engine,
1138 mut decode_handle,
1139 stream,
1140 _rate_watch,
1141 } = playback;
1142
1143 let _ = engine.stop();
1144 decode_handle.signal_stop();
1147 if let Some(stream) = stream {
1148 stream.abandon();
1149 }
1150 decode_handle.stop();
1151 drop(engine);
1152 }
1153
1154 fn stop_playback_and_clear_state(&mut self) {
1156 self.report(PlaybackReportState::Stopped);
1157 self.finish_play();
1158 self.stop_engine();
1159 self.timeline.reset();
1160 self.shared_state.set_playback_state(PlaybackState::Stopped);
1161 self.shared_state.set_position_ms(0);
1162 self.shared_state.set_track_info(None);
1163 }
1164
1165 pub fn remove_from_playlist(&mut self, id: QueueItemId) {
1172 let was_cursor = self.shared_state.is_cursor(id);
1173 let resume_after = was_cursor
1174 .then(|| self.shared_state.item_before(id))
1175 .flatten();
1176 self.shared_state.remove_item(id);
1177 if was_cursor {
1178 self.shared_state.set_cursor(resume_after);
1179 self.next_track();
1180 }
1181 }
1182
1183 pub fn track_ready(&mut self, id: QueueItemId) {
1186 self.shared_state.update_item_state(id, ItemState::Ready);
1188
1189 if !self.shared_state.is_cursor(id) {
1190 return;
1191 }
1192
1193 let is_playing = self.shared_state.playback_state() == PlaybackState::Playing;
1194 let current_track_id = self.shared_state.track_info().map(|t| t.id);
1195
1196 if is_playing && current_track_id == Some(id) {
1197 log::info!(
1200 "track_ready: download complete while streaming {:?}, refreshing metadata",
1201 id
1202 );
1203 self.refresh_track_metadata(id);
1204 return;
1205 }
1206
1207 if !is_playing && let Some(path) = self.shared_state.item_path_if_ready(id) {
1209 log::info!("track_ready: starting playback for {:?}", id);
1210 if let Err(e) = self.start_playback(id, &path, 0, Start::Playing) {
1211 log::error!("track_ready playback failed: {}", e);
1212 }
1213 }
1214 }
1215
1216 pub fn track_stream_ready(&mut self, id: QueueItemId) {
1219 if !self.shared_state.is_cursor(id) {
1220 return;
1221 }
1222
1223 let is_playing = self.shared_state.playback_state() == PlaybackState::Playing;
1224 if is_playing {
1225 return; }
1227
1228 match self.shared_state.item_playback_source(id) {
1229 Some(PlaybackSource::Streaming {
1230 path,
1231 bytes_written,
1232 total,
1233 }) => {
1234 log::info!("track_stream_ready: probing partial file for {:?}", id);
1235 self.probe_stream_for_playback(id, &path, bytes_written, total);
1236 }
1237 Some(PlaybackSource::Ready(path)) => {
1238 log::info!(
1240 "track_stream_ready: track already ready, starting normal playback for {:?}",
1241 id
1242 );
1243 if let Err(e) = self.start_playback(id, &path, 0, Start::Playing) {
1244 log::error!("track_stream_ready playback failed: {}", e);
1245 }
1246 }
1247 None => {} }
1249 }
1250
1251 fn refresh_track_metadata(&mut self, id: QueueItemId) {
1255 use crate::index::metadata;
1256
1257 let path = match self.shared_state.item_path_if_ready(id) {
1258 Some(p) => p,
1259 None => return,
1260 };
1261
1262 match metadata::read_metadata(&path) {
1263 Ok(meta) => {
1264 self.shared_state.update_item_metadata(
1265 id,
1266 meta.title,
1267 meta.artist,
1268 meta.album_artist.unwrap_or_default(),
1269 meta.album,
1270 meta.duration_ms.map(|d| d as u64),
1271 );
1272
1273 if let Some(current) = self.shared_state.track_info()
1282 && current.id == id
1283 {
1284 let probed = buffer::probe_file(&path).ok();
1285 let duration_ms = probed
1286 .as_ref()
1287 .map(|s| s.duration_ms)
1288 .filter(|d| *d > current.duration_ms)
1289 .unwrap_or(current.duration_ms);
1290 if duration_ms != current.duration_ms {
1291 log::info!(
1292 "track_ready: duration corrected {}ms → {}ms",
1293 current.duration_ms,
1294 duration_ms
1295 );
1296 }
1297 self.shared_state.set_track_info(Some(TrackInfo {
1298 duration_ms,
1299 path: path.clone(),
1300 ..current
1301 }));
1302 }
1303
1304 self.shared_state.signal_metadata_refresh();
1306 log::info!("track_ready: metadata refreshed for {:?}", id);
1307 }
1308 Err(e) => {
1309 log::warn!("track_ready: metadata refresh failed for {:?}: {}", id, e);
1310 }
1311 }
1312 }
1313
1314 fn on_track_changed(&mut self, id: QueueItemId, position_ms: u64) {
1322 if let Some(f) = self.in_flight.as_mut().filter(|f| f.item == id) {
1323 f.jump(position_ms);
1324 return;
1325 }
1326 self.finish_play();
1327 let track_id = self.shared_state.item_db_id(id);
1328 self.in_flight = Some(InFlight::new(id, track_id, position_ms));
1329 if let (Some(track_id), Some(recorder)) = (track_id, self.history.as_ref()) {
1330 recorder.record(PlayEvent::Started {
1331 track_id,
1332 position_ms,
1333 });
1334 }
1335 }
1336
1337 fn bank_listening(&mut self) {
1341 if let Some(f) = self.in_flight.as_mut()
1342 && let Some(at) = self.timeline.position_of(f.item)
1343 {
1344 f.advance(at);
1345 }
1346 }
1347
1348 fn report(&self, state: PlaybackReportState) {
1352 let Some(track_id) = self.in_flight.as_ref().and_then(InFlight::track_id) else {
1353 return;
1354 };
1355 if let Some(recorder) = self.history.as_ref() {
1356 recorder.record(PlayEvent::Playback(PlaybackReport {
1357 track_id,
1358 state,
1359 position_ms: self.shared_state.position_ms(),
1360 }));
1361 }
1362 }
1363
1364 fn finish_play(&mut self) -> Option<PlayEvent> {
1367 self.bank_listening();
1368 let flight = self.in_flight.take()?;
1369 let event = PlayEvent::Finished {
1370 track_id: flight.track_id()?,
1371 listened_ms: flight.listened_ms(),
1372 };
1373 if let Some(recorder) = self.history.as_ref() {
1374 recorder.record(event);
1375 }
1376 Some(event)
1377 }
1378
1379 pub fn update_playback_state(&mut self) {
1382 let Some(playback) = self.active_playback.as_ref() else {
1383 return;
1384 };
1385
1386 if self
1389 .lead_in_ends
1390 .is_some_and(|end| std::time::Instant::now() >= end)
1391 {
1392 self.lead_in_ends = None;
1393 self.shared_state.changed();
1394 }
1395
1396 if self.shared_state.playback_state() == PlaybackState::Paused
1397 && playback.engine.is_running()
1398 && playback.engine.is_silent()
1399 && let Err(e) = playback.engine.stop()
1400 {
1401 log::error!("stopping after fade failed: {}", e);
1402 }
1403
1404 let Some((id, position_ms)) = self.timeline.playhead() else {
1407 return;
1408 };
1409 if self.in_flight.as_ref().is_none_or(|f| f.item != id) {
1410 self.on_track_changed(id, position_ms);
1411 }
1412
1413 let current_id = self.shared_state.track_info().map(|t| t.id);
1416 if current_id != Some(id)
1417 && let Some((id, path, info, _)) = self.timeline.current_playback()
1418 {
1419 log::info!("timeline: now playing {:?}", id);
1420 self.shared_state.set_track_info(Some(TrackInfo {
1421 id,
1422 path,
1423 codec: info.codec,
1424 sample_rate: info.sample_rate,
1425 bit_depth: info.bit_depth,
1426 bitrate_kbps: info.bitrate_kbps,
1427 channels: info.channels,
1428 duration_ms: info.duration_ms,
1429 }));
1430 self.shared_state.set_cursor(Some(id));
1431 }
1432 }
1433
1434 pub fn track_failed(&mut self, id: QueueItemId) {
1441 if !self.shared_state.is_cursor(id) {
1442 return;
1443 }
1444 if self.shared_state.playback_state() != PlaybackState::Stopped {
1448 return;
1449 }
1450 log::info!("track {:?} cannot load, moving on", id);
1451 self.next_track();
1452 }
1453
1454 fn on_decode_finished(&mut self) {
1461 log::info!("decode finished, checking for next track");
1462 match self.shared_state.advance_cursor_loadable() {
1463 Some(id) => self.play(id),
1464 None => {
1465 log::info!("no more tracks — stopping");
1466 self.stop_playback_and_clear_state();
1467 }
1468 }
1469 }
1470
1471 fn snapshot_for_undo(
1475 &self,
1476 ids: &[QueueItemId],
1477 ) -> Vec<(Box<state::PlaylistItem>, Option<QueueItemId>)> {
1478 self.shared_state
1479 .items_before(ids)
1480 .into_iter()
1481 .filter_map(|(id, after)| Some((Box::new(self.shared_state.get_item(id)?), after)))
1482 .collect()
1483 }
1484
1485 fn push_undo(&mut self, entry: UndoEntry) {
1487 if let Some(ref mut batch) = self.batch_buffer {
1488 batch.push(entry);
1489 } else {
1490 self.undo_stack.push(entry);
1491 }
1492 }
1493
1494 pub fn process_command(&mut self, cmd: PlayerCommand) {
1496 match cmd {
1497 PlayerCommand::Play(id) => self.play(id),
1498 PlayerCommand::Cue {
1499 id,
1500 position_ms,
1501 play,
1502 } => self.cue(
1503 id,
1504 position_ms,
1505 if play { Start::Playing } else { Start::Paused },
1506 ),
1507 PlayerCommand::Pause => self.pause(),
1508 PlayerCommand::Resume => self.resume(),
1509 PlayerCommand::Stop => self.stop(),
1510 PlayerCommand::Seek(pos) => self.seek(pos),
1511 PlayerCommand::NextTrack => {
1512 let now = std::time::Instant::now();
1514 if now.duration_since(self.last_skip).as_millis() >= 150 {
1515 self.last_skip = now;
1516 self.next_track();
1517 }
1518 }
1519 PlayerCommand::PrevTrack => {
1520 let now = std::time::Instant::now();
1521 if now.duration_since(self.last_skip).as_millis() >= 150 {
1522 self.last_skip = now;
1523 self.prev_track();
1524 }
1525 }
1526 PlayerCommand::AddToPlaylist(items) => {
1527 let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
1528 self.shared_state.add_items(items);
1529 self.push_undo(UndoEntry::Added { ids });
1530 }
1531 PlayerCommand::UpdatePaths(updates) => {
1532 self.shared_state.update_paths(&updates);
1533 if let Some(info) = self.shared_state.track_info()
1534 && let Some((_, new_path)) = updates.iter().find(|(id, _)| *id == info.id)
1535 {
1536 self.shared_state.set_track_info(Some(TrackInfo {
1537 path: new_path.clone(),
1538 ..info
1539 }));
1540 }
1541 }
1542 PlayerCommand::InsertInPlaylist { items, after } => {
1543 let ids: Vec<QueueItemId> = items.iter().map(|i| i.id).collect();
1544 self.shared_state.insert_items_after(items, after);
1545 self.push_undo(UndoEntry::Inserted { ids });
1546 }
1547 PlayerCommand::ClearPlaylist => {
1548 self.stop_playback_and_clear_state();
1552 let (items, cursor) = self.shared_state.snapshot_playlist();
1553 self.shared_state.clear_playlist();
1554 self.push_undo(UndoEntry::Replaced { items, cursor });
1555 }
1556 PlayerCommand::ReplacePlaylist { items, start } => {
1557 self.stop_playback_and_clear_state();
1561 let (old_items, cursor) = self.shared_state.snapshot_playlist();
1562 self.shared_state.clear_playlist();
1563 self.push_undo(UndoEntry::Replaced {
1564 items: old_items,
1565 cursor,
1566 });
1567
1568 if items.is_empty() {
1569 return;
1570 }
1571 let start_id = items.get(start).unwrap_or(&items[0]).id;
1572 self.shared_state.add_items(items);
1573 self.play(start_id);
1574 }
1575 PlayerCommand::RemoveFromPlaylist(id) => {
1576 let item = self.shared_state.get_item(id);
1577 let after = self.shared_state.item_before(id);
1578 self.remove_from_playlist(id);
1579 if let Some(item) = item {
1580 self.push_undo(UndoEntry::Removed {
1581 items: vec![(Box::new(item), after)],
1582 });
1583 }
1584 }
1585 PlayerCommand::RemoveFromPlaylistBatch(ids) => {
1586 let items_with_pos = self.snapshot_for_undo(&ids);
1590 let resume_after = match self.shared_state.cursor() {
1591 Some(cursor) if ids.contains(&cursor) => {
1592 Some(self.shared_state.surviving_item_before(cursor, &ids))
1593 }
1594 _ => None,
1595 };
1596
1597 self.shared_state.remove_items(&ids);
1598
1599 if let Some(resume_after) = resume_after {
1600 self.shared_state.set_cursor(resume_after);
1601 self.next_track();
1602 }
1603
1604 if !items_with_pos.is_empty() {
1605 self.push_undo(UndoEntry::Removed {
1606 items: items_with_pos,
1607 });
1608 }
1609 }
1610 PlayerCommand::MoveInPlaylist { id, target, after } => {
1611 let was_after = self.shared_state.item_before(id);
1612 self.shared_state.move_item(id, target, after);
1613 self.push_undo(UndoEntry::Moved { id, was_after });
1614 }
1615 PlayerCommand::MoveItemsInPlaylist { ids, target, after } => {
1616 let entries = self.shared_state.items_before(&ids);
1617 self.shared_state.move_items(&ids, target, after);
1618 self.push_undo(UndoEntry::MovedBatch { entries });
1619 }
1620 PlayerCommand::ReorderPlaylist(order) => {
1621 let entries = self.shared_state.items_before(&order);
1625 self.shared_state.reorder_to(&order);
1626 self.push_undo(UndoEntry::MovedBatch { entries });
1627 }
1628 PlayerCommand::TrackReady(id) => self.track_ready(id),
1629 PlayerCommand::DecodeFinished => self.on_decode_finished(),
1630 PlayerCommand::TrackQueued => {}
1633 PlayerCommand::TrackStreamReady(id) => self.track_stream_ready(id),
1634 PlayerCommand::StreamProbed { id, info, mode } => self.stream_probed(id, *info, mode),
1635 PlayerCommand::TrackFailed(id) => self.track_failed(id),
1636 PlayerCommand::Undo => self.execute_undo(),
1637 PlayerCommand::Redo => self.execute_redo(),
1638 PlayerCommand::BeginUndoBatch => {
1639 self.batch_buffer = Some(Vec::new());
1640 }
1641 PlayerCommand::EndUndoBatch => {
1642 if let Some(entries) = self.batch_buffer.take() {
1643 if entries.len() == 1 {
1644 self.undo_stack.push(entries.into_iter().next().unwrap());
1646 } else if !entries.is_empty() {
1647 self.undo_stack.push(UndoEntry::Batch(entries));
1648 }
1649 }
1650 }
1651 PlayerCommand::SetOutputDevice(name) => self.set_output_device(name),
1652 PlayerCommand::RestartOutput => {
1653 log::info!("restarting audio output");
1654 self.restart_on_current_track();
1655 }
1656 PlayerCommand::ClearOutputDevice => self.clear_output_device(),
1657 }
1658 }
1659
1660 fn apply_entry(&mut self, entry: UndoEntry) -> UndoEntry {
1662 match entry {
1663 UndoEntry::Added { ids } | UndoEntry::Inserted { ids } => {
1664 let items_with_pos = self.snapshot_for_undo(&ids);
1666 self.shared_state.remove_items(&ids);
1667 UndoEntry::Removed {
1668 items: items_with_pos,
1669 }
1670 }
1671 UndoEntry::Removed { items } => {
1672 let mut ids = Vec::with_capacity(items.len());
1674 for (item, after) in items {
1675 ids.push(item.id);
1676 self.shared_state.insert_item_at(*item, after);
1677 }
1678 UndoEntry::Added { ids }
1679 }
1680 UndoEntry::Moved { id, was_after } => {
1681 let current_after = self.shared_state.item_before(id);
1682 self.shared_state.move_item_to(id, was_after);
1683 UndoEntry::Moved {
1684 id,
1685 was_after: current_after,
1686 }
1687 }
1688 UndoEntry::MovedBatch { entries } => {
1689 let ids: Vec<QueueItemId> = entries.iter().map(|(id, _)| *id).collect();
1690 let current_positions = self.shared_state.items_before(&ids);
1691 self.shared_state.move_items_to(&entries);
1692 UndoEntry::MovedBatch {
1693 entries: current_positions,
1694 }
1695 }
1696 UndoEntry::Replaced { items, cursor } => {
1697 let (current_items, current_cursor) = self.shared_state.snapshot_playlist();
1698 self.shared_state.restore_playlist(items, cursor);
1699 UndoEntry::Replaced {
1700 items: current_items,
1701 cursor: current_cursor,
1702 }
1703 }
1704 UndoEntry::Batch(entries) => {
1705 let mut inverses: Vec<_> = entries
1707 .into_iter()
1708 .rev()
1709 .map(|e| self.apply_entry(e))
1710 .collect();
1711 inverses.reverse();
1712 UndoEntry::Batch(inverses)
1713 }
1714 }
1715 }
1716
1717 fn reconcile_playback(&mut self) {
1731 let Some(playing) = self.shared_state.track_info().map(|t| t.id) else {
1732 return;
1733 };
1734 if self.shared_state.get_item(playing).is_some() {
1735 return;
1736 }
1737 let resume = (self.shared_state.playback_state() == PlaybackState::Playing)
1742 .then(|| self.shared_state.cursor())
1743 .flatten();
1744 self.stop_playback_and_clear_state();
1745 if let Some(id) = resume {
1746 self.play(id);
1747 }
1748 }
1749
1750 fn execute_undo(&mut self) {
1752 let Some(entry) = self.undo_stack.pop_undo() else {
1753 return;
1754 };
1755 let inverse = self.apply_entry(entry);
1756 self.undo_stack.push_redo(inverse);
1757 self.reconcile_playback();
1758 }
1759
1760 fn execute_redo(&mut self) {
1762 let Some(entry) = self.undo_stack.pop_redo() else {
1763 return;
1764 };
1765 let inverse = self.apply_entry(entry);
1766 self.undo_stack.push_undo_keep_redo(inverse);
1767 self.reconcile_playback();
1768 }
1769
1770 pub fn run(&mut self) {
1776 use crossbeam_channel::RecvTimeoutError;
1777
1778 let rx = self.commands.rx.clone();
1779 loop {
1780 let received = match self.next_wake() {
1781 Some(at) => rx.recv_deadline(at),
1782 None => rx.recv().map_err(|_| RecvTimeoutError::Disconnected),
1783 };
1784 match received {
1785 Ok(cmd) => self.process_command(cmd),
1786 Err(RecvTimeoutError::Timeout) => {}
1787 Err(RecvTimeoutError::Disconnected) => break,
1788 }
1789 self.update_playback_state();
1790 }
1791 self.stop();
1792 }
1793
1794 fn next_wake(&self) -> Option<std::time::Instant> {
1798 let playback = self.active_playback.as_ref()?;
1799 let now = std::time::Instant::now();
1800 match self.shared_state.playback_state() {
1801 PlaybackState::Playing => {
1802 let next_track = self
1803 .timeline
1804 .until_next_track()
1805 .map(|left| now + left + BOUNDARY_SLACK);
1806 match (next_track, self.lead_in_ends) {
1807 (Some(a), Some(b)) => Some(a.min(b)),
1808 (a, b) => a.or(b),
1809 }
1810 }
1811 PlaybackState::Paused if playback.engine.is_running() => Some(now + FADE_CHECK),
1812 _ => None,
1813 }
1814 }
1815
1816 pub fn spawn() -> (
1819 Arc<SharedPlayerState>,
1820 Arc<PlaybackTimeline>,
1821 Arc<VizSnapshot>,
1822 crossbeam_channel::Sender<PlayerCommand>,
1823 ) {
1824 let mut player = Self::new();
1825 player.history = PlayRecorder::spawn();
1826 let state = player.shared_state();
1827 let timeline = player.timeline();
1828 let viz_snapshot = player.viz_snapshot();
1829 let tx = player.command_sender();
1830
1831 thread::Builder::new()
1832 .name("koan-player".into())
1833 .spawn(move || player.run())
1834 .expect("failed to spawn player thread");
1835
1836 (state, timeline, viz_snapshot, tx)
1837 }
1838}
1839
1840#[cfg(test)]
1841mod tests {
1842 #[test]
1843 fn a_download_in_progress_is_known_by_its_own_extension() {
1844 use std::path::Path;
1845 assert_eq!(
1847 lengthless_mode_for(Path::new("/c/t.m4a.part")),
1848 streaming::ProbeMode::LengthlessWholeEnd
1849 );
1850 assert_eq!(
1851 lengthless_mode_for(Path::new("/c/t.OPUS.part")),
1852 streaming::ProbeMode::LengthlessWholeEnd
1853 );
1854 assert_eq!(
1855 lengthless_mode_for(Path::new("/c/t.flac.part")),
1856 streaming::ProbeMode::Lengthless
1857 );
1858 assert_eq!(
1859 media_extension(Path::new("/c/t.m4a.part")).as_deref(),
1860 Some("m4a")
1861 );
1862 assert_eq!(
1863 media_extension(Path::new("/c/t.mp3")).as_deref(),
1864 Some("mp3")
1865 );
1866 }
1867
1868 use super::*;
1869 use state::PlaylistItem;
1870 use std::path::PathBuf;
1871 use std::sync::atomic::{AtomicU64, Ordering};
1872
1873 fn make_item(title: &str) -> PlaylistItem {
1874 PlaylistItem {
1875 playlist_entry_id: None,
1876 id: QueueItemId::new(),
1877 db_id: None,
1878 path: PathBuf::from(format!("/music/{title}.flac")),
1879 title: title.to_string(),
1880 artist: String::new(),
1881 album_artist: String::new(),
1882 album: String::new(),
1883 year: None,
1884 codec: None,
1885 track_number: None,
1886 disc: None,
1887 duration_ms: None,
1888 state: ItemState::Ready,
1889 }
1890 }
1891
1892 fn playlist_ids(player: &Player) -> Vec<QueueItemId> {
1893 let (items, _) = player.shared_state.snapshot_playlist();
1894 items.iter().map(|i| i.id).collect()
1895 }
1896
1897 fn playlist_titles(player: &Player) -> Vec<String> {
1898 let (items, _) = player.shared_state.snapshot_playlist();
1899 items.iter().map(|i| i.title.clone()).collect()
1900 }
1901
1902 fn pending_item(title: &str) -> PlaylistItem {
1903 PlaylistItem {
1904 playlist_entry_id: None,
1905 state: ItemState::Pending,
1906 ..make_item(title)
1907 }
1908 }
1909
1910 fn pretend_playing(player: &mut Player, id: QueueItemId) {
1914 let item = player
1915 .shared_state
1916 .get_item(id)
1917 .expect("item is in the queue");
1918 player.shared_state.set_track_info(Some(TrackInfo {
1919 id,
1920 path: item.path,
1921 codec: String::new(),
1922 sample_rate: 44_100,
1923 bit_depth: None,
1924 bitrate_kbps: None,
1925 channels: 2,
1926 duration_ms: 1_000,
1927 }));
1928 player
1929 .shared_state
1930 .set_playback_state(PlaybackState::Playing);
1931 }
1932
1933 fn playing_id(player: &Player) -> Option<QueueItemId> {
1934 player.shared_state.track_info().map(|t| t.id)
1935 }
1936
1937 fn seed(player: &mut Player, n: usize) -> Vec<QueueItemId> {
1939 let items: Vec<_> = (0..n).map(|i| make_item(&format!("t{i}"))).collect();
1940 let ids = items.iter().map(|i| i.id).collect();
1941 player.process_command(PlayerCommand::AddToPlaylist(items));
1942 ids
1943 }
1944
1945 fn listen(player: &mut Player, from_ms: u64, to_ms: u64) {
1949 let mut at = from_ms;
1950 if let Some(f) = player.in_flight.as_mut() {
1951 f.advance(at); }
1953 while at < to_ms {
1954 at = (at + 50).min(to_ms);
1955 if let Some(f) = player.in_flight.as_mut() {
1956 f.advance(at);
1957 }
1958 }
1959 }
1960
1961 fn start(player: &mut Player, track_id: i64) -> QueueItemId {
1962 let id = QueueItemId::new();
1963 player.on_track_changed(id, 0);
1964 player
1966 .in_flight
1967 .as_mut()
1968 .unwrap()
1969 .track_id_for_test(track_id);
1970 id
1971 }
1972
1973 #[test]
1974 fn a_gapless_transition_closes_the_outgoing_track_and_opens_the_next() {
1975 let mut player = Player::new();
1976 start(&mut player, 11);
1977 listen(&mut player, 0, 200_000);
1978
1979 let b = QueueItemId::new();
1980 player.on_track_changed(b, 0);
1981 let f = player
1982 .in_flight
1983 .as_ref()
1984 .expect("the next track is counting");
1985 assert_eq!(f.item, b);
1986 assert_eq!(f.listened_ms(), 0, "and starts from nothing");
1987 }
1988
1989 #[test]
1990 fn a_track_skipped_seconds_in_is_still_history() {
1991 let mut player = Player::new();
1992 start(&mut player, 7);
1993 listen(&mut player, 0, 2_000);
1994
1995 let event = player
1996 .finish_play()
1997 .expect("putting something on is a thing you did, however briefly");
1998 assert!(matches!(
1999 event,
2000 history::PlayEvent::Finished {
2001 track_id: 7,
2002 listened_ms: 2_000
2003 }
2004 ));
2005 }
2006
2007 #[test]
2008 fn a_track_is_closed_out_once() {
2009 let mut player = Player::new();
2010 start(&mut player, 7);
2011 listen(&mut player, 0, 200_000);
2012
2013 assert!(player.finish_play().is_some());
2014 assert!(player.finish_play().is_none());
2015 }
2016
2017 #[test]
2018 fn seeking_around_a_track_does_not_enter_it_twice() {
2019 let mut player = Player::new();
2020 let id = start(&mut player, 7);
2021 listen(&mut player, 0, 120_000);
2022
2023 player.on_track_changed(id, 30_000);
2025 assert_eq!(
2026 player.in_flight.as_ref().unwrap().listened_ms(),
2027 120_000,
2028 "the seek kept the count rather than restarting it"
2029 );
2030 listen(&mut player, 30_000, 40_000);
2031
2032 let Some(history::PlayEvent::Finished { listened_ms, .. }) = player.finish_play() else {
2033 panic!("still one play");
2034 };
2035 assert_eq!(listened_ms, 130_000);
2036 assert!(player.finish_play().is_none());
2037 }
2038
2039 #[test]
2040 fn a_track_that_is_not_in_the_library_is_not_recorded() {
2041 let mut player = Player::new();
2042 let id = QueueItemId::new();
2043 player.on_track_changed(id, 0);
2044 listen(&mut player, 0, 200_000);
2045 assert!(player.finish_play().is_none());
2046 }
2047
2048 #[test]
2049 fn stopping_closes_out_what_was_heard() {
2050 let mut player = Player::new();
2051 start(&mut player, 7);
2052 listen(&mut player, 0, 150_000);
2053
2054 player.stop_playback_and_clear_state();
2055 assert!(player.in_flight.is_none(), "the stop consumed it");
2056 }
2057
2058 #[test]
2059 fn resume_with_nothing_loaded_plays_the_cursor() {
2060 let dir = tempfile::tempdir().unwrap();
2061 let path = dir.path().join("t.wav");
2062 crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
2063
2064 let mut player = Player::new();
2065 player.backend = Box::new(StuckBackend {
2066 rate: 8_000.0,
2067 asked: Default::default(),
2068 starts: Default::default(),
2069 });
2070 let item = PlaylistItem {
2071 db_id: Some(5),
2072 path,
2073 ..make_item("t")
2074 };
2075 let id = item.id;
2076 player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
2077 player.shared_state.set_cursor(Some(id));
2078 assert!(player.active_playback.is_none());
2079
2080 player.process_command(PlayerCommand::Resume);
2081 assert!(
2082 player.active_playback.is_some(),
2083 "the cursor's track started"
2084 );
2085 assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
2086 player.process_command(PlayerCommand::Stop);
2087 }
2088
2089 #[test]
2090 fn a_player_with_nothing_coming_has_nothing_to_wake_for() {
2091 let dir = tempfile::tempdir().unwrap();
2092 let path = dir.path().join("t.wav");
2093 crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
2094
2095 let mut player = Player::new();
2096 player.backend = Box::new(StuckBackend {
2097 rate: 8_000.0,
2098 asked: Default::default(),
2099 starts: Default::default(),
2100 });
2101 let item = PlaylistItem {
2102 path,
2103 ..make_item("t")
2104 };
2105 let id = item.id;
2106 player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
2107 assert_eq!(player.next_wake(), None, "stopped");
2108
2109 player.process_command(PlayerCommand::Play(id));
2110 assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
2111 assert_eq!(
2112 player.next_wake(),
2113 None,
2114 "playing, with no track queued after it"
2115 );
2116
2117 player.pause_now();
2118 assert_eq!(player.next_wake(), None, "paused");
2119 player.process_command(PlayerCommand::Stop);
2120 }
2121
2122 #[test]
2123 fn a_track_cued_paused_never_starts_the_output() {
2124 let dir = tempfile::tempdir().unwrap();
2125 let path = dir.path().join("t.wav");
2126 crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
2127
2128 let starts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
2129 let mut player = Player::new();
2130 player.backend = Box::new(StuckBackend {
2131 rate: 8_000.0,
2132 asked: Default::default(),
2133 starts: starts.clone(),
2134 });
2135 let item = PlaylistItem {
2136 path,
2137 ..make_item("t")
2138 };
2139 let id = item.id;
2140 player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
2141
2142 player.process_command(PlayerCommand::Cue {
2143 id,
2144 position_ms: 2_000,
2145 play: false,
2146 });
2147 assert!(player.active_playback.is_some(), "loaded");
2148 assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
2149 loop {
2153 match player
2154 .commands
2155 .rx
2156 .recv_timeout(std::time::Duration::from_secs(5))
2157 {
2158 Ok(PlayerCommand::TrackQueued) => break,
2159 Ok(_) => {}
2160 Err(e) => panic!("the decoder never queued the track: {e}"),
2161 }
2162 }
2163 let at = player.shared_state.position_ms();
2164 assert!((1_750..=2_000).contains(&at), "cued at {at}ms");
2165 assert_eq!(starts.load(Ordering::Relaxed), 0, "not a sample let out");
2166
2167 player.process_command(PlayerCommand::Seek(4_000));
2169 assert_eq!(player.shared_state.playback_state(), PlaybackState::Paused);
2170 assert_eq!(starts.load(Ordering::Relaxed), 0);
2171
2172 player.process_command(PlayerCommand::Resume);
2173 assert_eq!(player.shared_state.playback_state(), PlaybackState::Playing);
2174 assert_eq!(starts.load(Ordering::Relaxed), 1);
2175 player.process_command(PlayerCommand::Stop);
2176 }
2177
2178 #[test]
2179 fn the_server_hears_each_turn_playback_takes() {
2180 use PlaybackReportState::{Paused, Playing, Stopped};
2181 use history::PlaybackReport;
2182
2183 let dir = tempfile::tempdir().unwrap();
2184 let path = dir.path().join("t.wav");
2185 crate::test_utils::generate_wav(&path, 8_000, 1, 10.0, 16);
2186
2187 let mut player = Player::new();
2188 player.backend = Box::new(StuckBackend {
2189 rate: 8_000.0,
2190 asked: Default::default(),
2191 starts: Default::default(),
2192 });
2193 let (recorder, events) = PlayRecorder::capture();
2194 player.history = Some(recorder);
2195
2196 let item = PlaylistItem {
2197 db_id: Some(5),
2198 path,
2199 ..make_item("t")
2200 };
2201 let id = item.id;
2202 player.process_command(PlayerCommand::AddToPlaylist(vec![item]));
2203 player.process_command(PlayerCommand::Play(id));
2204 player.process_command(PlayerCommand::Pause);
2205 player.process_command(PlayerCommand::Seek(4_000));
2206 player.process_command(PlayerCommand::Resume);
2207 player.process_command(PlayerCommand::Stop);
2208
2209 let report = |state, position_ms| {
2210 PlayEvent::Playback(PlaybackReport {
2211 track_id: 5,
2212 state,
2213 position_ms,
2214 })
2215 };
2216 assert_eq!(
2217 events.try_iter().collect::<Vec<_>>(),
2218 vec![
2219 PlayEvent::Started {
2220 track_id: 5,
2221 position_ms: 0
2222 },
2223 report(Paused, 0),
2224 report(Paused, 4_000),
2225 report(Playing, 4_000),
2226 report(Stopped, 4_000),
2227 PlayEvent::Finished {
2228 track_id: 5,
2229 listened_ms: 0
2230 },
2231 ]
2232 );
2233 }
2234
2235 #[test]
2236 fn removing_the_playing_track_resumes_at_its_successor() {
2237 let mut player = Player::new();
2238 let ids = seed(&mut player, 5);
2239 player.shared_state.set_cursor(Some(ids[2]));
2240
2241 player.process_command(PlayerCommand::RemoveFromPlaylist(ids[2]));
2242
2243 assert_eq!(
2244 player.shared_state.cursor(),
2245 Some(ids[3]),
2246 "playback must continue at the next track, not restart the queue"
2247 );
2248 assert_eq!(player.playback_starts, 1);
2249 }
2250
2251 #[test]
2252 fn removing_the_first_playing_track_resumes_at_the_new_first() {
2253 let mut player = Player::new();
2254 let ids = seed(&mut player, 3);
2255 player.shared_state.set_cursor(Some(ids[0]));
2256
2257 player.process_command(PlayerCommand::RemoveFromPlaylist(ids[0]));
2258
2259 assert_eq!(player.shared_state.cursor(), Some(ids[1]));
2260 }
2261
2262 #[test]
2263 fn next_track_parks_on_a_track_that_has_not_downloaded_yet() {
2264 let mut player = Player::new();
2265 let playing = make_item("playing");
2266 let waiting = pending_item("waiting");
2267 let later = make_item("later");
2268 let (playing_id, waiting_id) = (playing.id, waiting.id);
2269 player.process_command(PlayerCommand::AddToPlaylist(vec![playing, waiting, later]));
2270 player.shared_state.set_cursor(Some(playing_id));
2271
2272 player.process_command(PlayerCommand::DecodeFinished);
2273
2274 assert_eq!(
2275 player.shared_state.cursor(),
2276 Some(waiting_id),
2277 "the cursor parks on the track being fetched"
2278 );
2279 assert_eq!(
2280 player.playback_starts, 0,
2281 "nothing to play until its bytes land"
2282 );
2283
2284 player
2287 .shared_state
2288 .update_item_state(waiting_id, ItemState::Ready);
2289 player.process_command(PlayerCommand::TrackReady(waiting_id));
2290
2291 assert_eq!(player.playback_starts, 1);
2292 assert_eq!(player.shared_state.cursor(), Some(waiting_id));
2293 }
2294
2295 #[test]
2296 fn a_download_that_cannot_land_moves_the_cursor_on() {
2297 let mut player = Player::new();
2298 let waiting = pending_item("waiting");
2299 let later = make_item("later");
2300 let (waiting_id, later_id) = (waiting.id, later.id);
2301 player.process_command(PlayerCommand::AddToPlaylist(vec![waiting, later]));
2302
2303 player.process_command(PlayerCommand::Play(waiting_id));
2304 assert_eq!(player.playback_starts, 0, "nothing to play yet");
2305
2306 player
2308 .shared_state
2309 .update_item_state(waiting_id, ItemState::Failed("remote unavailable".into()));
2310 player.process_command(PlayerCommand::TrackFailed(waiting_id));
2311
2312 assert_eq!(
2313 player.shared_state.cursor(),
2314 Some(later_id),
2315 "the queue moves past a track that can never load"
2316 );
2317 assert_eq!(player.playback_starts, 1);
2318 }
2319
2320 #[test]
2321 fn a_queue_that_can_never_load_stops_rather_than_waiting() {
2322 let mut player = Player::new();
2323 let first = pending_item("first");
2324 let second = pending_item("second");
2325 let (first_id, second_id) = (first.id, second.id);
2326 player.process_command(PlayerCommand::AddToPlaylist(vec![first, second]));
2327
2328 player.process_command(PlayerCommand::Play(first_id));
2329 for id in [first_id, second_id] {
2330 player
2331 .shared_state
2332 .update_item_state(id, ItemState::Failed("remote unavailable".into()));
2333 player.process_command(PlayerCommand::TrackFailed(id));
2334 }
2335
2336 assert_eq!(player.playback_starts, 0);
2337 assert_eq!(
2338 player.shared_state.playback_state(),
2339 PlaybackState::Stopped,
2340 "a stop the UI can see, not an indefinite wait for TrackReady"
2341 );
2342 }
2343
2344 #[test]
2345 fn a_failure_elsewhere_in_the_queue_leaves_the_cursor_alone() {
2346 let mut player = Player::new();
2347 let waiting = pending_item("waiting");
2348 let other = pending_item("other");
2349 let (waiting_id, other_id) = (waiting.id, other.id);
2350 player.process_command(PlayerCommand::AddToPlaylist(vec![waiting, other]));
2351 player.process_command(PlayerCommand::Play(waiting_id));
2352
2353 player
2354 .shared_state
2355 .update_item_state(other_id, ItemState::Failed("remote unavailable".into()));
2356 player.process_command(PlayerCommand::TrackFailed(other_id));
2357
2358 assert_eq!(
2359 player.shared_state.cursor(),
2360 Some(waiting_id),
2361 "a track still downloading keeps the cursor"
2362 );
2363 }
2364
2365 #[test]
2366 fn batch_delete_containing_the_cursor_restarts_the_engine_once() {
2367 let mut player = Player::new();
2368 let ids = seed(&mut player, 5);
2369 player.shared_state.set_cursor(Some(ids[2]));
2370
2371 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![
2372 ids[1], ids[2], ids[3],
2373 ]));
2374
2375 assert_eq!(playlist_titles(&player), vec!["t0", "t4"]);
2376 assert_eq!(player.shared_state.cursor(), Some(ids[4]));
2377 assert_eq!(
2378 player.playback_starts, 1,
2379 "one resume for the whole selection, not one per deleted track"
2380 );
2381 }
2382
2383 #[test]
2384 fn batch_delete_below_the_cursor_leaves_playback_alone() {
2385 let mut player = Player::new();
2386 let ids = seed(&mut player, 4);
2387 player.shared_state.set_cursor(Some(ids[0]));
2388
2389 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![ids[2], ids[3]]));
2390
2391 assert_eq!(player.shared_state.cursor(), Some(ids[0]));
2392 assert_eq!(player.playback_starts, 0);
2393 }
2394
2395 #[test]
2396 fn undo_of_a_batch_delete_restores_the_original_order() {
2397 let mut player = Player::new();
2401 let items = vec![
2402 make_item("A"),
2403 make_item("B"),
2404 make_item("C"),
2405 make_item("D"),
2406 ];
2407 let (b_id, c_id) = (items[1].id, items[2].id);
2408 player.process_command(PlayerCommand::AddToPlaylist(items));
2409
2410 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![c_id, b_id]));
2411 assert_eq!(playlist_titles(&player), vec!["A", "D"]);
2412
2413 player.process_command(PlayerCommand::Undo);
2414 assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
2415 }
2416
2417 #[test]
2420 fn undo_add_removes_items() {
2421 let mut player = Player::new();
2422 let items = vec![make_item("A"), make_item("B")];
2423 let ids: Vec<_> = items.iter().map(|i| i.id).collect();
2424
2425 player.process_command(PlayerCommand::AddToPlaylist(items));
2426 assert_eq!(playlist_ids(&player), ids);
2427 assert!(player.undo_stack().can_undo());
2428
2429 player.process_command(PlayerCommand::Undo);
2430 assert!(playlist_ids(&player).is_empty());
2431 assert!(player.undo_stack().can_redo());
2432 }
2433
2434 #[test]
2435 fn redo_add_restores_items() {
2436 let mut player = Player::new();
2437 let items = vec![make_item("A"), make_item("B")];
2438
2439 player.process_command(PlayerCommand::AddToPlaylist(items));
2440 player.process_command(PlayerCommand::Undo);
2441 assert!(playlist_ids(&player).is_empty());
2442
2443 player.process_command(PlayerCommand::Redo);
2444 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2445 }
2446
2447 #[test]
2450 fn undo_remove_restores_item_at_position() {
2451 let mut player = Player::new();
2452 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2453 let b_id = items[1].id;
2454
2455 player.process_command(PlayerCommand::AddToPlaylist(items));
2456 player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
2457 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2458
2459 player.process_command(PlayerCommand::Undo);
2460 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2461 }
2462
2463 #[test]
2464 fn undo_remove_first_item() {
2465 let mut player = Player::new();
2466 let items = vec![make_item("A"), make_item("B")];
2467 let a_id = items[0].id;
2468
2469 player.process_command(PlayerCommand::AddToPlaylist(items));
2470 player.process_command(PlayerCommand::RemoveFromPlaylist(a_id));
2471 assert_eq!(playlist_titles(&player), vec!["B"]);
2472
2473 player.process_command(PlayerCommand::Undo);
2474 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2475 }
2476
2477 #[test]
2478 fn undo_batch_remove_restores_all() {
2479 let mut player = Player::new();
2480 let items = vec![
2481 make_item("A"),
2482 make_item("B"),
2483 make_item("C"),
2484 make_item("D"),
2485 ];
2486 let b_id = items[1].id;
2487 let c_id = items[2].id;
2488
2489 player.process_command(PlayerCommand::AddToPlaylist(items));
2490 let version_before = player.shared_state.playlist_version();
2491 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![b_id, c_id]));
2492 assert_eq!(playlist_titles(&player), vec!["A", "D"]);
2493 assert_eq!(
2496 player.shared_state.playlist_version(),
2497 version_before + 1,
2498 "batch removal must bump the playlist version exactly once"
2499 );
2500
2501 player.process_command(PlayerCommand::Undo);
2503 assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
2504 }
2505
2506 #[test]
2507 fn redo_batch_remove() {
2508 let mut player = Player::new();
2509 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2510 let a_id = items[0].id;
2511 let b_id = items[1].id;
2512
2513 player.process_command(PlayerCommand::AddToPlaylist(items));
2514 player.process_command(PlayerCommand::RemoveFromPlaylistBatch(vec![a_id, b_id]));
2515 player.process_command(PlayerCommand::Undo);
2516 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2517
2518 player.process_command(PlayerCommand::Redo);
2519 assert_eq!(playlist_titles(&player), vec!["C"]);
2520 }
2521
2522 #[test]
2523 fn redo_remove() {
2524 let mut player = Player::new();
2525 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2526 let b_id = items[1].id;
2527
2528 player.process_command(PlayerCommand::AddToPlaylist(items));
2529 player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
2530 player.process_command(PlayerCommand::Undo);
2531 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2532
2533 player.process_command(PlayerCommand::Redo);
2534 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2535 }
2536
2537 #[test]
2540 fn undo_insert_removes_inserted_items() {
2541 let mut player = Player::new();
2542 let items = vec![make_item("A"), make_item("C")];
2543 let a_id = items[0].id;
2544
2545 player.process_command(PlayerCommand::AddToPlaylist(items));
2546
2547 let inserted = vec![make_item("B")];
2548 player.process_command(PlayerCommand::InsertInPlaylist {
2549 items: inserted,
2550 after: a_id,
2551 });
2552 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2553
2554 player.process_command(PlayerCommand::Undo);
2555 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2556 }
2557
2558 #[test]
2561 fn undo_move_restores_position() {
2562 let mut player = Player::new();
2563 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2564 let a_id = items[0].id;
2565 let c_id = items[2].id;
2566
2567 player.process_command(PlayerCommand::AddToPlaylist(items));
2568
2569 player.process_command(PlayerCommand::MoveInPlaylist {
2571 id: a_id,
2572 target: c_id,
2573 after: true,
2574 });
2575 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2576
2577 player.process_command(PlayerCommand::Undo);
2578 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2579 }
2580
2581 #[test]
2582 fn redo_move() {
2583 let mut player = Player::new();
2584 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2585 let a_id = items[0].id;
2586 let c_id = items[2].id;
2587
2588 player.process_command(PlayerCommand::AddToPlaylist(items));
2589 player.process_command(PlayerCommand::MoveInPlaylist {
2590 id: a_id,
2591 target: c_id,
2592 after: true,
2593 });
2594 player.process_command(PlayerCommand::Undo);
2595 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2596
2597 player.process_command(PlayerCommand::Redo);
2598 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2599 }
2600
2601 #[test]
2604 fn undo_batch_move() {
2605 let mut player = Player::new();
2606 let items = vec![
2607 make_item("A"),
2608 make_item("B"),
2609 make_item("C"),
2610 make_item("D"),
2611 ];
2612 let a_id = items[0].id;
2613 let b_id = items[1].id;
2614 let d_id = items[3].id;
2615
2616 player.process_command(PlayerCommand::AddToPlaylist(items));
2617
2618 player.process_command(PlayerCommand::MoveItemsInPlaylist {
2620 ids: vec![a_id, b_id],
2621 target: d_id,
2622 after: true,
2623 });
2624 assert_eq!(playlist_titles(&player), vec!["C", "D", "A", "B"]);
2625
2626 player.process_command(PlayerCommand::Undo);
2627 assert_eq!(playlist_titles(&player), vec!["A", "B", "C", "D"]);
2628 }
2629
2630 #[test]
2633 fn undo_clear_restores_playlist() {
2634 let mut player = Player::new();
2635 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2636
2637 player.process_command(PlayerCommand::AddToPlaylist(items));
2638 player.process_command(PlayerCommand::ClearPlaylist);
2639 assert!(playlist_ids(&player).is_empty());
2640
2641 player.process_command(PlayerCommand::Undo);
2642 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2643 }
2644
2645 #[test]
2650 fn undoing_a_replace_does_not_leave_the_engine_on_an_orphaned_track() {
2651 let mut player = Player::new();
2652 let original = seed(&mut player, 3);
2653 player.shared_state.set_cursor(Some(original[0]));
2654 pretend_playing(&mut player, original[0]);
2655
2656 let replacement = vec![make_item("something else")];
2657 let orphan = replacement[0].id;
2658 player.process_command(PlayerCommand::ReplacePlaylist {
2659 items: replacement,
2660 start: 0,
2661 });
2662 pretend_playing(&mut player, orphan);
2664
2665 player.process_command(PlayerCommand::Undo);
2666
2667 assert_eq!(playlist_ids(&player), original, "the queue comes back");
2668 assert!(
2669 player.shared_state.get_item(orphan).is_none(),
2670 "and the replacement is gone from it"
2671 );
2672 assert!(
2673 playing_id(&player).is_none_or(|id| player.shared_state.get_item(id).is_some()),
2674 "so nothing may still be playing out of it"
2675 );
2676 }
2677
2678 #[test]
2680 fn undoing_an_add_does_not_leave_the_engine_on_a_removed_track() {
2681 let mut player = Player::new();
2682 seed(&mut player, 2);
2683 let added = seed(&mut player, 1);
2684 pretend_playing(&mut player, added[0]);
2685
2686 player.process_command(PlayerCommand::Undo);
2687
2688 assert!(player.shared_state.get_item(added[0]).is_none());
2689 assert!(
2690 playing_id(&player).is_none_or(|id| player.shared_state.get_item(id).is_some()),
2691 "the engine cannot be left on the item the undo removed"
2692 );
2693 }
2694
2695 #[test]
2697 fn undoing_a_move_leaves_playback_alone() {
2698 let mut player = Player::new();
2699 let ids = seed(&mut player, 3);
2700 player.shared_state.set_cursor(Some(ids[0]));
2701 pretend_playing(&mut player, ids[0]);
2702 let starts = player.playback_starts;
2703
2704 player.process_command(PlayerCommand::MoveInPlaylist {
2705 id: ids[2],
2706 target: ids[0],
2707 after: false,
2708 });
2709 player.process_command(PlayerCommand::Undo);
2710
2711 assert_eq!(playlist_ids(&player), ids);
2712 assert_eq!(playing_id(&player), Some(ids[0]), "still on the same track");
2713 assert_eq!(player.playback_starts, starts, "and not restarted");
2714 }
2715
2716 #[test]
2717 fn redo_clear() {
2718 let mut player = Player::new();
2719 let items = vec![make_item("A"), make_item("B")];
2720
2721 player.process_command(PlayerCommand::AddToPlaylist(items));
2722 player.process_command(PlayerCommand::ClearPlaylist);
2723 player.process_command(PlayerCommand::Undo);
2724 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2725
2726 player.process_command(PlayerCommand::Redo);
2727 assert!(playlist_ids(&player).is_empty());
2728 }
2729
2730 #[test]
2733 fn multiple_undos_in_sequence() {
2734 let mut player = Player::new();
2735
2736 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("A")]));
2737 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
2738 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("C")]));
2739 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2740
2741 player.process_command(PlayerCommand::Undo);
2742 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2743
2744 player.process_command(PlayerCommand::Undo);
2745 assert_eq!(playlist_titles(&player), vec!["A"]);
2746
2747 player.process_command(PlayerCommand::Undo);
2748 assert!(playlist_ids(&player).is_empty());
2749 }
2750
2751 #[test]
2752 fn undo_redo_undo_cycle() {
2753 let mut player = Player::new();
2754 let items = vec![make_item("A"), make_item("B")];
2755
2756 player.process_command(PlayerCommand::AddToPlaylist(items));
2757 player.process_command(PlayerCommand::Undo);
2758 assert!(playlist_ids(&player).is_empty());
2759
2760 player.process_command(PlayerCommand::Redo);
2761 assert_eq!(playlist_titles(&player), vec!["A", "B"]);
2762
2763 player.process_command(PlayerCommand::Undo);
2764 assert!(playlist_ids(&player).is_empty());
2765 }
2766
2767 #[test]
2768 fn new_action_clears_redo_stack() {
2769 let mut player = Player::new();
2770 let items = vec![make_item("A")];
2771
2772 player.process_command(PlayerCommand::AddToPlaylist(items));
2773 player.process_command(PlayerCommand::Undo);
2774 assert!(player.undo_stack().can_redo());
2775
2776 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("B")]));
2778 assert!(!player.undo_stack().can_redo());
2779 }
2780
2781 #[test]
2782 fn undo_on_empty_stack_is_noop() {
2783 let mut player = Player::new();
2784 player.process_command(PlayerCommand::Undo);
2785 assert!(playlist_ids(&player).is_empty());
2786 }
2787
2788 #[test]
2789 fn redo_on_empty_stack_is_noop() {
2790 let mut player = Player::new();
2791 player.process_command(PlayerCommand::Redo);
2792 assert!(playlist_ids(&player).is_empty());
2793 }
2794
2795 #[test]
2798 fn playback_commands_not_undoable() {
2799 let mut player = Player::new();
2800 player.process_command(PlayerCommand::Pause);
2801 player.process_command(PlayerCommand::Resume);
2802 player.process_command(PlayerCommand::NextTrack);
2803 player.process_command(PlayerCommand::PrevTrack);
2804 assert!(!player.undo_stack().can_undo());
2805 }
2806
2807 #[test]
2808 fn update_paths_not_undoable() {
2809 let mut player = Player::new();
2810 let items = vec![make_item("A")];
2811 let id = items[0].id;
2812 player.process_command(PlayerCommand::AddToPlaylist(items));
2813
2814 let undo_count = player.undo_stack().undo_len();
2815 player.process_command(PlayerCommand::UpdatePaths(vec![(
2816 id,
2817 PathBuf::from("/new/path.flac"),
2818 )]));
2819 assert_eq!(player.undo_stack().undo_len(), undo_count);
2820 }
2821
2822 #[test]
2825 fn add_remove_undo_undo_produces_original() {
2826 let mut player = Player::new();
2827 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2828 let b_id = items[1].id;
2829 let original_titles = vec!["A", "B", "C"];
2830
2831 player.process_command(PlayerCommand::AddToPlaylist(items));
2832 player.process_command(PlayerCommand::RemoveFromPlaylist(b_id));
2833 assert_eq!(playlist_titles(&player), vec!["A", "C"]);
2834
2835 player.process_command(PlayerCommand::Undo);
2837 assert_eq!(playlist_titles(&player), original_titles);
2838
2839 player.process_command(PlayerCommand::Undo);
2841 assert!(playlist_ids(&player).is_empty());
2842 }
2843
2844 #[test]
2845 fn interleaved_adds_and_moves_undo() {
2846 let mut player = Player::new();
2847 let items = vec![make_item("A"), make_item("B"), make_item("C")];
2848 let a_id = items[0].id;
2849 let c_id = items[2].id;
2850
2851 player.process_command(PlayerCommand::AddToPlaylist(items));
2852
2853 player.process_command(PlayerCommand::MoveInPlaylist {
2855 id: a_id,
2856 target: c_id,
2857 after: true,
2858 });
2859 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2860
2861 player.process_command(PlayerCommand::AddToPlaylist(vec![make_item("D")]));
2863 assert_eq!(playlist_titles(&player), vec!["B", "C", "A", "D"]);
2864
2865 player.process_command(PlayerCommand::Undo);
2867 assert_eq!(playlist_titles(&player), vec!["B", "C", "A"]);
2868
2869 player.process_command(PlayerCommand::Undo);
2871 assert_eq!(playlist_titles(&player), vec!["A", "B", "C"]);
2872 }
2873
2874 #[test]
2879 fn stop_engine_drops_engine_synchronously() {
2880 use std::sync::atomic::{AtomicBool, Ordering};
2881
2882 struct MockEngine {
2883 dropped: Arc<AtomicBool>,
2884 }
2885 impl AudioEngineHandle for MockEngine {
2886 fn start(&self) -> Result<(), BackendError> {
2887 Ok(())
2888 }
2889 fn stop(&self) -> Result<(), BackendError> {
2890 Ok(())
2891 }
2892 fn is_running(&self) -> bool {
2893 false
2894 }
2895 fn fade_out(&self) {}
2896 fn fade_in(&self) -> Result<(), BackendError> {
2897 Ok(())
2898 }
2899 fn is_silent(&self) -> bool {
2900 false
2901 }
2902 }
2903 impl Drop for MockEngine {
2904 fn drop(&mut self) {
2905 self.dropped.store(true, Ordering::SeqCst);
2906 }
2907 }
2908
2909 let dropped = Arc::new(AtomicBool::new(false));
2910
2911 let stop_flag = Arc::new(AtomicBool::new(false));
2913 let decode_handle = buffer::DecodeHandle::new_for_test(stop_flag);
2914
2915 let mut player = Player::new();
2916 player.active_playback = Some(ActivePlayback {
2917 engine: Box::new(MockEngine {
2918 dropped: dropped.clone(),
2919 }),
2920 decode_handle,
2921 stream: None,
2922 _rate_watch: None,
2923 });
2924
2925 player.stop_engine();
2926
2927 assert!(
2931 dropped.load(Ordering::SeqCst),
2932 "AudioEngine must be dropped synchronously in stop_engine (GitHub #89)"
2933 );
2934 }
2935
2936 #[test]
2937 fn abandoning_a_stream_wakes_a_reader_parked_at_the_write_head() {
2938 let live = LiveStream {
2939 feed: crate::remote::downloads::ByteFeed::new(),
2940 abandoned: Default::default(),
2941 };
2942 let feed = live.feed.clone();
2943 let started = std::time::Instant::now();
2944 let reader = thread::spawn(move || {
2945 feed.wait_past(
2946 0,
2947 std::time::Instant::now() + std::time::Duration::from_secs(30),
2948 )
2949 });
2950 thread::sleep(std::time::Duration::from_millis(50));
2951 live.abandon();
2952 reader.join().unwrap();
2953
2954 assert!(live.abandoned.load(std::sync::atomic::Ordering::Acquire));
2955 assert!(started.elapsed() < std::time::Duration::from_secs(5));
2956 }
2957
2958 struct StuckBackend {
2963 rate: f64,
2964 asked: Arc<std::sync::Mutex<Option<(f64, u32)>>>,
2965 starts: Arc<std::sync::atomic::AtomicUsize>,
2967 }
2968
2969 struct NullEngine {
2970 starts: Arc<std::sync::atomic::AtomicUsize>,
2971 running: std::sync::atomic::AtomicBool,
2972 lead_in: Arc<AtomicU64>,
2973 }
2974 impl AudioEngineHandle for NullEngine {
2975 fn start(&self) -> Result<(), BackendError> {
2976 self.starts.fetch_add(1, Ordering::Relaxed);
2977 self.running.store(true, Ordering::Relaxed);
2978 Ok(())
2979 }
2980 fn stop(&self) -> Result<(), BackendError> {
2981 self.running.store(false, Ordering::Relaxed);
2982 Ok(())
2983 }
2984 fn is_running(&self) -> bool {
2985 self.running.load(Ordering::Relaxed)
2986 }
2987 fn fade_out(&self) {}
2988 fn fade_in(&self) -> Result<(), BackendError> {
2989 Ok(())
2990 }
2991 fn is_silent(&self) -> bool {
2992 false
2993 }
2994 fn lead_in(&self, frames: u64) {
2995 self.lead_in.store(frames, Ordering::Relaxed);
2996 }
2997 }
2998
2999 impl AudioBackend for StuckBackend {
3000 fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
3001 Ok(vec![self.default_device()?])
3002 }
3003 fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
3004 Ok(backend::DeviceInfo {
3005 name: "Stuck DAC".into(),
3006 sample_rates: vec![self.rate],
3007 platform_id: 0,
3008 })
3009 }
3010 fn supported_sample_rates(
3011 &self,
3012 _device: &backend::DeviceInfo,
3013 ) -> Result<Vec<f64>, BackendError> {
3014 Ok(vec![self.rate])
3015 }
3016 fn get_device_sample_rate(
3017 &self,
3018 _device: &backend::DeviceInfo,
3019 ) -> Result<f64, BackendError> {
3020 Ok(self.rate)
3021 }
3022 fn set_device_sample_rate(
3023 &self,
3024 _device: &backend::DeviceInfo,
3025 rate: f64,
3026 ) -> Result<f64, BackendError> {
3027 Err(BackendError::UnsupportedSampleRate(rate))
3028 }
3029 fn create_engine(
3030 &self,
3031 _device: &backend::DeviceInfo,
3032 sample_rate: f64,
3033 channels: u32,
3034 _consumer: rtrb::Consumer<f32>,
3035 _samples_played: Arc<AtomicU64>,
3036 ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
3037 *self.asked.lock().unwrap() = Some((sample_rate, channels));
3038 Ok(Box::new(NullEngine {
3039 starts: self.starts.clone(),
3040 running: Default::default(),
3041 lead_in: Default::default(),
3042 }))
3043 }
3044 }
3045
3046 fn engine_format_for(source_rate: u32, channels: u16, device_rate: f64) -> (f64, u32) {
3047 let asked = Arc::new(std::sync::Mutex::new(None));
3048 let mut player = Player::new();
3049 player.backend = Box::new(StuckBackend {
3050 rate: device_rate,
3051 asked: asked.clone(),
3052 starts: Default::default(),
3053 });
3054
3055 let info = buffer::StreamInfo {
3056 codec: "MP3".into(),
3057 sample_rate: source_rate,
3058 channels,
3059 bit_depth: Some(16),
3060 bitrate_kbps: None,
3061 duration_ms: 1000,
3062 };
3063 let (_producer, consumer) = rtrb::RingBuffer::new(16);
3064 player
3065 .create_engine_for(&info, consumer)
3066 .expect("engine creation should succeed");
3067 let asked = *asked.lock().unwrap();
3068 asked.expect("engine was never created")
3069 }
3070
3071 #[test]
3072 fn engine_uses_source_rate_when_device_refuses_switch() {
3073 assert_eq!(engine_format_for(22050, 2, 48000.0), (22050.0, 2));
3076 assert_eq!(engine_format_for(32000, 2, 44100.0), (32000.0, 2));
3077 }
3078
3079 #[test]
3080 fn engine_uses_source_channel_count() {
3081 assert_eq!(engine_format_for(44100, 1, 44100.0), (44100.0, 1));
3082 }
3083
3084 fn settled_rate_for(source_rate: u32, device_rate: f64) -> Option<u32> {
3086 let mut player = Player::new();
3087 player.backend = Box::new(StuckBackend {
3088 rate: device_rate,
3089 asked: Arc::new(std::sync::Mutex::new(None)),
3090 starts: Default::default(),
3091 });
3092 let state = player.shared_state.clone();
3093
3094 let info = buffer::StreamInfo {
3095 codec: "MP3".into(),
3096 sample_rate: source_rate,
3097 channels: 2,
3098 bit_depth: Some(16),
3099 bitrate_kbps: None,
3100 duration_ms: 1000,
3101 };
3102 let (_producer, consumer) = rtrb::RingBuffer::new(16);
3103 player
3104 .create_engine_for(&info, consumer)
3105 .expect("engine creation should succeed");
3106 state.output_sample_rate()
3107 }
3108
3109 #[test]
3110 fn settled_device_rate_reaches_the_shared_state() {
3111 assert_eq!(settled_rate_for(22050, 48000.0), Some(48000));
3115 assert_eq!(settled_rate_for(44100, 44100.0), Some(44100));
3117 }
3118
3119 struct SlowBackend {
3121 observed: Arc<std::sync::Mutex<Vec<Option<u32>>>>,
3122 state: Arc<SharedPlayerState>,
3123 lead_in: Arc<AtomicU64>,
3124 }
3125
3126 impl AudioBackend for SlowBackend {
3127 fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
3128 Ok(vec![self.default_device()?])
3129 }
3130 fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
3131 Ok(backend::DeviceInfo {
3132 name: "Slow DAC".into(),
3133 sample_rates: vec![44100.0, 48000.0],
3134 platform_id: 0,
3135 })
3136 }
3137 fn supported_sample_rates(
3138 &self,
3139 _device: &backend::DeviceInfo,
3140 ) -> Result<Vec<f64>, BackendError> {
3141 Ok(vec![44100.0, 48000.0])
3142 }
3143 fn get_device_sample_rate(
3144 &self,
3145 _device: &backend::DeviceInfo,
3146 ) -> Result<f64, BackendError> {
3147 Ok(48000.0)
3148 }
3149 fn set_device_sample_rate(
3150 &self,
3151 _device: &backend::DeviceInfo,
3152 rate: f64,
3153 ) -> Result<f64, BackendError> {
3154 self.observed
3156 .lock()
3157 .unwrap()
3158 .push(self.state.output_sample_rate());
3159 Ok(rate)
3160 }
3161 fn create_engine(
3162 &self,
3163 _device: &backend::DeviceInfo,
3164 _sample_rate: f64,
3165 _channels: u32,
3166 _consumer: rtrb::Consumer<f32>,
3167 _samples_played: Arc<AtomicU64>,
3168 ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
3169 Ok(Box::new(NullEngine {
3170 starts: Default::default(),
3171 running: Default::default(),
3172 lead_in: self.lead_in.clone(),
3173 }))
3174 }
3175 }
3176
3177 #[test]
3178 fn the_previous_rate_is_not_published_while_the_device_reclocks() {
3179 let mut player = Player::new();
3185 let state = player.shared_state.clone();
3186 state.set_output_sample_rate(48000);
3187
3188 let observed = Arc::new(std::sync::Mutex::new(Vec::new()));
3189 player.backend = Box::new(SlowBackend {
3190 observed: observed.clone(),
3191 state: state.clone(),
3192 lead_in: Default::default(),
3193 });
3194
3195 let info = buffer::StreamInfo {
3196 codec: "FLAC".into(),
3197 sample_rate: 44100,
3198 channels: 2,
3199 bit_depth: Some(16),
3200 bitrate_kbps: None,
3201 duration_ms: 1000,
3202 };
3203 let (_producer, consumer) = rtrb::RingBuffer::new(16);
3204 player
3205 .create_engine_for(&info, consumer)
3206 .expect("engine creation should succeed");
3207
3208 assert_eq!(
3209 *observed.lock().unwrap(),
3210 vec![None],
3211 "mid-switch the output rate must read as unknown, not as the last track's"
3212 );
3213 assert_eq!(state.output_sample_rate(), Some(44100));
3214 }
3215
3216 fn lead_in_for(source_rate: u32) -> u64 {
3219 let mut player = Player::new();
3220 let lead_in = Arc::new(AtomicU64::new(0));
3221 player.backend = Box::new(SlowBackend {
3222 observed: Default::default(),
3223 state: player.shared_state.clone(),
3224 lead_in: lead_in.clone(),
3225 });
3226 let info = buffer::StreamInfo {
3227 codec: "FLAC".into(),
3228 sample_rate: source_rate,
3229 channels: 2,
3230 bit_depth: Some(16),
3231 bitrate_kbps: None,
3232 duration_ms: 1000,
3233 };
3234 let (_producer, consumer) = rtrb::RingBuffer::new(16);
3235 player
3236 .create_engine_for(&info, consumer)
3237 .expect("engine creation should succeed");
3238 lead_in.load(Ordering::Relaxed)
3239 }
3240
3241 #[test]
3242 fn an_engine_made_inside_the_silence_keeps_the_rest_of_it() {
3243 let mut player = Player::new();
3246 let lead_in = Arc::new(AtomicU64::new(0));
3247 player.backend = Box::new(SlowBackend {
3248 observed: Default::default(),
3249 state: player.shared_state.clone(),
3250 lead_in: lead_in.clone(),
3251 });
3252 let info = |sample_rate| buffer::StreamInfo {
3253 codec: "FLAC".into(),
3254 sample_rate,
3255 channels: 2,
3256 bit_depth: Some(16),
3257 bitrate_kbps: None,
3258 duration_ms: 1000,
3259 };
3260 let (_p, consumer) = rtrb::RingBuffer::new(16);
3261 player.create_engine_for(&info(44100), consumer).unwrap();
3262 let first = lead_in.swap(0, Ordering::Relaxed);
3263 assert!(first > 0);
3264
3265 let (_p, consumer) = rtrb::RingBuffer::new(16);
3266 player.create_engine_for(&info(48000), consumer).unwrap();
3267 let carried = lead_in.load(Ordering::Relaxed);
3268 assert!(
3270 carried > 0 && carried < first * 48000 / 44100,
3271 "carried {carried} of {first}"
3272 );
3273 }
3274
3275 #[test]
3276 fn silence_leads_in_only_after_the_device_changed_rate() {
3277 assert!(
3278 lead_in_for(44100) > 0,
3279 "the device is relocking, so the start of the track would be lost"
3280 );
3281 assert_eq!(lead_in_for(48000), 0, "no switch, nothing to wait for");
3282 }
3283
3284 struct WatchedBackend {
3286 inner: StuckBackend,
3287 #[allow(clippy::type_complexity)]
3288 captured: Arc<std::sync::Mutex<Option<Box<dyn Fn(f64) + Send + Sync>>>>,
3289 }
3290
3291 struct NullWatch;
3292 impl backend::SampleRateWatch for NullWatch {}
3293
3294 impl AudioBackend for WatchedBackend {
3295 fn list_devices(&self) -> Result<Vec<backend::DeviceInfo>, BackendError> {
3296 self.inner.list_devices()
3297 }
3298 fn default_device(&self) -> Result<backend::DeviceInfo, BackendError> {
3299 self.inner.default_device()
3300 }
3301 fn supported_sample_rates(
3302 &self,
3303 device: &backend::DeviceInfo,
3304 ) -> Result<Vec<f64>, BackendError> {
3305 self.inner.supported_sample_rates(device)
3306 }
3307 fn get_device_sample_rate(
3308 &self,
3309 device: &backend::DeviceInfo,
3310 ) -> Result<f64, BackendError> {
3311 self.inner.get_device_sample_rate(device)
3312 }
3313 fn set_device_sample_rate(
3314 &self,
3315 device: &backend::DeviceInfo,
3316 rate: f64,
3317 ) -> Result<f64, BackendError> {
3318 self.inner.set_device_sample_rate(device, rate)
3319 }
3320 fn watch_device_sample_rate(
3321 &self,
3322 _device: &backend::DeviceInfo,
3323 on_change: Box<dyn Fn(f64) + Send + Sync>,
3324 ) -> Option<Box<dyn backend::SampleRateWatch>> {
3325 *self.captured.lock().unwrap() = Some(on_change);
3326 Some(Box::new(NullWatch))
3327 }
3328 fn create_engine(
3329 &self,
3330 device: &backend::DeviceInfo,
3331 sample_rate: f64,
3332 channels: u32,
3333 consumer: rtrb::Consumer<f32>,
3334 samples_played: Arc<AtomicU64>,
3335 ) -> Result<Box<dyn AudioEngineHandle>, BackendError> {
3336 self.inner
3337 .create_engine(device, sample_rate, channels, consumer, samples_played)
3338 }
3339 }
3340
3341 #[test]
3342 fn external_rate_change_reaches_the_shared_state() {
3343 let captured = Arc::new(std::sync::Mutex::new(None));
3347 let mut player = Player::new();
3348 player.backend = Box::new(WatchedBackend {
3349 inner: StuckBackend {
3350 rate: 44100.0,
3351 asked: Arc::new(std::sync::Mutex::new(None)),
3352 starts: Default::default(),
3353 },
3354 captured: captured.clone(),
3355 });
3356 let state = player.shared_state.clone();
3357
3358 let info = buffer::StreamInfo {
3359 codec: "FLAC".into(),
3360 sample_rate: 44100,
3361 channels: 2,
3362 bit_depth: Some(16),
3363 bitrate_kbps: None,
3364 duration_ms: 1000,
3365 };
3366 let (_producer, consumer) = rtrb::RingBuffer::new(16);
3367 player
3368 .create_engine_for(&info, consumer)
3369 .expect("engine creation should succeed");
3370 assert_eq!(state.output_sample_rate(), Some(44100));
3371
3372 let on_change = captured.lock().unwrap().take().expect("watch registered");
3373 on_change(48000.0);
3374 assert_eq!(state.output_sample_rate(), Some(48000));
3375 }
3376}