1use std::collections::{HashSet, VecDeque};
2use std::panic::AssertUnwindSafe;
3use std::sync::Arc;
4use std::time::Duration;
5
6use parking_lot::{Condvar, Mutex};
7
8use crate::config;
9use crate::db::queries;
10use crate::helpers::download_track;
11use crate::player::commands::PlayerCommand;
12use crate::player::state::{ItemState, QueueItemId, SharedPlayerState};
13use crate::remote::downloads::{self, DownloadStore};
14
15const PRIORITY_PERMITS: usize = 2;
19
20const CANCEL_CHECK: Duration = Duration::from_millis(250);
22
23#[derive(Clone)]
40pub struct DownloadQueue {
41 inner: Arc<Inner>,
42}
43
44struct Inner {
45 queue: Mutex<Queue>,
50 has_work: Condvar,
51 state: Arc<SharedPlayerState>,
52 cmd_tx: crossbeam_channel::Sender<PlayerCommand>,
53 last_evicted: Mutex<Option<std::time::Instant>>,
55 spawned: std::sync::atomic::AtomicUsize,
58}
59
60fn workers_allowed() -> usize {
63 config::Config::cached()
64 .remote
65 .download_workers
66 .clamp(1, 16)
67}
68
69fn ensure_workers(inner: &Arc<Inner>) {
71 use std::sync::atomic::Ordering;
72 let want = workers_allowed();
73 loop {
74 let have = inner.spawned.load(Ordering::Relaxed);
75 if have >= want {
76 return;
77 }
78 if inner
79 .spawned
80 .compare_exchange(have, have + 1, Ordering::Relaxed, Ordering::Relaxed)
81 .is_err()
82 {
83 continue;
84 }
85 let worker = inner.clone();
86 if let Err(e) = std::thread::Builder::new()
87 .name(format!("koan-dl-{have}"))
88 .spawn(move || {
89 if have == 0 {
93 trim_cache(&worker, false);
94 }
95 worker_loop(worker, have)
96 })
97 {
98 log::error!("failed to spawn download worker {have}: {e}");
99 inner.spawned.fetch_sub(1, Ordering::Relaxed);
100 return;
101 }
102 }
103}
104
105fn next_item(
111 q: &mut Queue,
112 store: &DownloadStore,
113 cursor: Option<(i64, QueueItemId)>,
114) -> Option<Job> {
115 let entry = match cursor {
116 Some((db_id, _)) if store.in_flight(db_id) => return None,
117 Some((_, queue_id)) => q
118 .pending
119 .iter()
120 .position(|(_, qid)| *qid == queue_id)
121 .and_then(|ix| q.pending.remove(ix))
122 .or_else(|| q.pending.pop_front()),
123 None => q.pending.pop_front(),
124 };
125 entry
126 .map(|(db_id, id)| (db_id, Some(id)))
127 .or_else(|| q.cache.pop_front().map(|db_id| (db_id, None)))
128}
129
130fn retry_server_now() {
135 if let Some(client) = crate::helpers::subsonic_client(&config::Config::cached()) {
136 client.outage().retry_now();
137 }
138}
139
140fn cursor_download(state: &SharedPlayerState) -> Option<(i64, QueueItemId)> {
142 let id = state.cursor()?;
143 let item = state.get_item(id)?;
144 if item.state != ItemState::Pending {
145 return None;
146 }
147 Some((item.db_id?, id))
148}
149
150type Job = (i64, Option<QueueItemId>);
153
154#[derive(Default)]
157struct Queue {
158 pending: VecDeque<(i64, QueueItemId)>,
160 cache: VecDeque<i64>,
162 priority_active: usize,
163 window: Option<HashSet<i64>>,
167}
168
169#[derive(Debug, PartialEq, Eq)]
171enum Dispatch {
172 Spawn,
174 Requeued,
176 AlreadyRunning,
178}
179
180fn claim_priority(q: &mut Queue, store: &DownloadStore, item: (i64, QueueItemId)) -> Dispatch {
184 let (db_id, queue_id) = item;
185 q.pending.retain(|(_, qid)| *qid != queue_id);
186
187 if store.join(db_id, Some(queue_id)) {
190 return Dispatch::AlreadyRunning;
191 }
192 if q.priority_active >= PRIORITY_PERMITS {
193 q.pending.push_front(item);
194 return Dispatch::Requeued;
195 }
196 q.priority_active += 1;
197 store.claim(db_id, Some(queue_id));
198 Dispatch::Spawn
199}
200
201struct Permit {
203 inner: Arc<Inner>,
204}
205
206impl Drop for Permit {
207 fn drop(&mut self) {
208 let mut q = self.inner.queue.lock();
209 q.priority_active = q.priority_active.saturating_sub(1);
210 drop(q);
211 self.inner.has_work.notify_all();
212 }
213}
214
215impl DownloadQueue {
216 pub fn spawn(
219 cmd_tx: crossbeam_channel::Sender<PlayerCommand>,
220 state: Arc<SharedPlayerState>,
221 ) -> Self {
222 let inner = Arc::new(Inner {
223 queue: Mutex::new(Queue::default()),
224 has_work: Condvar::new(),
225 state,
226 cmd_tx,
227 last_evicted: Mutex::new(None),
228 spawned: std::sync::atomic::AtomicUsize::new(0),
229 });
230
231 let watcher_inner = inner.clone();
232 if let Err(e) = std::thread::Builder::new()
233 .name("koan-dl-watch".into())
234 .spawn(move || follow_playlist(watcher_inner))
235 {
236 log::error!("failed to spawn download playlist watcher: {}", e);
237 }
238
239 Self { inner }
240 }
241
242 pub fn cache(&self, track_ids: Vec<i64>) {
245 if track_ids.is_empty() {
246 return;
247 }
248 self.inner.queue.lock().cache.extend(track_ids);
249 wake_workers(&self.inner);
250 }
251}
252
253fn wake_workers(inner: &Arc<Inner>) {
256 ensure_workers(inner);
257 retry_server_now();
258 inner.has_work.notify_all();
259}
260
261fn sync(inner: &Arc<Inner>) {
280 let window = playback_window(&inner.state);
281 let mut q = inner.queue.lock();
282 let wanted = inner.state.pending_downloads();
283 let entries = window
284 .as_ref()
285 .map(|w| w.iter().map(|(_, id)| *id).collect::<HashSet<_>>());
286 let added = sync_with(&mut q, inner.state.downloads(), &wanted, entries.as_ref());
287 q.window = window.map(|w| w.into_iter().map(|(track, _)| track).collect());
288 drop(q);
289 if added {
290 wake_workers(inner);
291 }
292}
293
294fn playback_window(state: &SharedPlayerState) -> Option<Vec<(i64, QueueItemId)>> {
298 let limit = config::Config::cached().cache_limit_bytes()?;
299 let mut order = state.playback_order();
300 let tracks: Vec<i64> = order.iter().map(|(track, _)| *track).collect();
301 let fits = crate::db::pool::shared()
302 .get()
303 .map_err(|e| e.to_string())
304 .and_then(|db| {
305 crate::helpers::playback_window(&db, limit, &tracks).map_err(|e| e.to_string())
306 });
307 match fits {
308 Ok(fits) => {
309 order.truncate(fits);
310 Some(order)
311 }
312 Err(e) => {
313 log::warn!("cache window: {e}");
314 None
315 }
316 }
317}
318
319fn sync_with(
323 q: &mut Queue,
324 store: &DownloadStore,
325 wanted: &[(i64, QueueItemId)],
326 window: Option<&HashSet<QueueItemId>>,
327) -> bool {
328 let wanted: Vec<_> = match window {
329 Some(window) => wanted
330 .iter()
331 .filter(|(_, id)| window.contains(id))
332 .copied()
333 .collect(),
334 None => wanted.to_vec(),
335 };
336 let unfetched = store.resync(&wanted);
337 let before: HashSet<QueueItemId> = q.pending.iter().map(|(_, id)| *id).collect();
338 let added = unfetched.iter().any(|(_, id)| !before.contains(id));
339 q.pending = unfetched.into();
340 added
341}
342
343fn dispatch_priority(inner: &Arc<Inner>, item: (i64, QueueItemId)) {
345 let dispatch = claim_priority(&mut inner.queue.lock(), inner.state.downloads(), item);
346 match dispatch {
347 Dispatch::AlreadyRunning => {}
348 Dispatch::Requeued => {
349 inner.has_work.notify_one();
350 }
351 Dispatch::Spawn => {
352 let spawn_inner = inner.clone();
353 let spawned = std::thread::Builder::new()
354 .name("koan-dl-prio".into())
355 .spawn(move || {
356 let _permit = Permit {
357 inner: spawn_inner.clone(),
358 };
359 run_download(&spawn_inner, item.0);
360 });
361 if let Err(e) = spawned {
362 log::error!("failed to spawn priority download: {}", e);
363 let mut q = inner.queue.lock();
364 let store = inner.state.downloads();
365 for id in downloads::withdraw(store, item.0) {
366 q.pending.push_front((item.0, id));
367 }
368 q.priority_active = q.priority_active.saturating_sub(1);
369 drop(q);
370 inner.has_work.notify_one();
371 }
372 }
373 }
374}
375
376fn run_download(inner: &Arc<Inner>, db_id: i64) {
384 let store = inner.state.downloads();
385 let cfg = config::Config::cached();
386 let result = match crate::helpers::subsonic_client(&cfg) {
387 None => Some(Err(crate::helpers::remote_unavailable(&cfg))),
390 Some(client) => {
391 let checked = std::cell::Cell::new(std::time::Instant::now());
396 let cancelled = || {
397 if checked.get().elapsed() < CANCEL_CHECK {
398 return false;
399 }
400 checked.set(std::time::Instant::now());
401 store.abandoned(db_id)
402 };
403 std::panic::catch_unwind(AssertUnwindSafe(|| {
404 download_track(
405 db_id,
406 &cancelled,
407 &inner.cmd_tx,
408 &inner.state,
409 &cfg,
410 &client,
411 )
412 }))
413 .unwrap_or_else(|_| {
414 log::error!("download panicked for track {db_id}");
415 Some(Err("download panicked".into()))
416 })
417 }
418 };
419
420 if let Some(Ok(_)) = &result
424 && store.kept(db_id)
425 {
426 pin(db_id);
427 }
428 let mut q = inner.queue.lock();
429 let settled = match result {
430 Some(result) => Some(downloads::settle(&inner.state, db_id, &result)),
431 None => {
434 for id in downloads::withdraw(store, db_id) {
435 q.pending.push_front((db_id, id));
436 }
437 None
438 }
439 };
440 drop(q);
441 if let Some(settled) = settled {
442 settled.announce(&inner.cmd_tx);
443 }
444 inner.has_work.notify_all();
446 if config::Config::cached().cache_limit_bytes().is_some() {
448 sync(inner);
449 }
450 trim_cache(inner, false);
451}
452
453fn pin(db_id: i64) {
456 let pinned = crate::db::pool::shared()
457 .get()
458 .map_err(|e| e.to_string())
459 .and_then(|db| queries::pin_cached(&db.conn, &[db_id]).map_err(|e| e.to_string()));
460 if let Err(e) = pinned {
461 log::warn!("could not pin the download of track {db_id}: {e}");
462 }
463}
464
465const EVICT_EVERY: std::time::Duration = std::time::Duration::from_secs(60);
468
469fn trim_cache(inner: &Inner, now: bool) {
475 {
476 let mut last = inner.last_evicted.lock();
477 if !now && last.is_some_and(|t| t.elapsed() < EVICT_EVERY) {
478 return;
479 }
480 *last = Some(std::time::Instant::now());
481 }
482 let cfg = config::Config::cached();
483 if cfg.cache_limit_bytes().is_none() {
484 return;
485 }
486 let keep = inner.queue.lock().window.clone().unwrap_or_else(|| {
487 inner
488 .state
489 .playback_order()
490 .into_iter()
491 .map(|(track, _)| track)
492 .collect()
493 });
494 match crate::db::pool::shared().get() {
495 Ok(db) => {
496 if crate::helpers::evict_cache(&db, &cfg, &keep, false) > 0 {
497 inner.state.reset_items_with_missing_files();
500 }
501 }
502 Err(e) => log::warn!("cache eviction: could not open the database: {e}"),
503 }
504}
505
506static LIMIT_CHANGES: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
508
509pub fn cache_limit_changed() {
512 LIMIT_CHANGES.fetch_add(1, std::sync::atomic::Ordering::AcqRel);
513 crate::signal::engine_changed().bump();
514}
515
516fn worker_loop(inner: Arc<Inner>, index: usize) {
529 loop {
530 let item = loop {
531 if let Some(client) = crate::helpers::subsonic_client(&config::Config::cached()) {
532 client.outage().hold();
533 }
534 let cursor = cursor_download(&inner.state);
535 let mut q = inner.queue.lock();
536 if index >= workers_allowed() {
537 inner.has_work.wait(&mut q);
538 continue;
539 }
540 let store = inner.state.downloads();
541 match next_item(&mut q, store, cursor) {
542 Some((db_id, entry)) => {
545 if store.claim(db_id, entry) {
546 break db_id;
547 }
548 }
549 None => inner.has_work.wait(&mut q),
550 }
551 };
552 run_download(&inner, item);
553 }
554}
555
556fn follow_playlist(inner: Arc<Inner>) {
569 let changed = crate::signal::engine_changed();
570 let mut seen = changed.generation();
571 let mut last_version: Option<u64> = None;
572 let mut last_cursor: Option<QueueItemId> = None;
573 let mut last_limit = LIMIT_CHANGES.load(std::sync::atomic::Ordering::Acquire);
574 loop {
575 let version = inner.state.pending_version();
579 let current = inner.state.cursor();
580 let limit = LIMIT_CHANGES.load(std::sync::atomic::Ordering::Acquire);
581 if last_version != Some(version) || current != last_cursor || limit != last_limit {
582 last_version = Some(version);
583 sync(&inner);
584 }
585 if limit != last_limit {
586 last_limit = limit;
587 trim_cache(&inner, true);
588 }
589 if current != last_cursor {
590 last_cursor = current;
591 inner.has_work.notify_all();
592 if let Some(cursor_id) = current {
593 promote_cursor(&inner, cursor_id);
594 }
595 }
596 seen = changed.wait(seen);
597 }
598}
599
600fn promote_cursor(inner: &Arc<Inner>, cursor_id: QueueItemId) {
604 if inner.state.item_state(cursor_id) != Some(ItemState::Pending) {
605 return;
606 }
607 let mut priority_items = Vec::new();
608 {
609 let mut q = inner.queue.lock();
610 if let Some(pos) = q.pending.iter().position(|(_, qid)| *qid == cursor_id) {
611 priority_items.push(q.pending.remove(pos).expect("position just found"));
612 if let Some(next) = q.pending.pop_front() {
614 priority_items.push(next);
615 }
616 }
617 }
618 for item in priority_items {
619 dispatch_priority(inner, item);
620 }
621}
622
623#[cfg(test)]
624mod tests {
625 use super::*;
626 use crate::player::state::PlaylistItem;
627
628 fn qid() -> QueueItemId {
629 QueueItemId::new()
630 }
631
632 fn waiters(store: &DownloadStore, db_id: i64) -> HashSet<QueueItemId> {
633 store.waiters(db_id).into_iter().collect()
634 }
635
636 #[test]
637 fn priority_lane_never_exceeds_its_permits() {
638 let (mut q, store) = (Queue::default(), DownloadStore::new());
639
640 let mut spawned = 0;
642 for i in 0..500 {
643 if claim_priority(&mut q, &store, (i, qid())) == Dispatch::Spawn {
644 spawned += 1;
645 }
646 assert!(
647 q.priority_active <= PRIORITY_PERMITS,
648 "priority lane over its permit count at iteration {}",
649 i
650 );
651 }
652
653 assert_eq!(spawned, PRIORITY_PERMITS, "only permitted claims may spawn");
654 assert_eq!(
655 q.pending.len(),
656 500 - PRIORITY_PERMITS,
657 "everything else must be queued, not dropped"
658 );
659 }
660
661 #[test]
662 fn released_permits_are_reusable() {
663 let (mut q, store) = (Queue::default(), DownloadStore::new());
664 assert_eq!(claim_priority(&mut q, &store, (1, qid())), Dispatch::Spawn);
665 assert_eq!(claim_priority(&mut q, &store, (2, qid())), Dispatch::Spawn);
666 assert_eq!(
667 claim_priority(&mut q, &store, (3, qid())),
668 Dispatch::Requeued
669 );
670
671 let _ = downloads::withdraw(&store, 1);
672 q.priority_active -= 1;
673 assert_eq!(claim_priority(&mut q, &store, (4, qid())), Dispatch::Spawn);
674 assert!(q.priority_active <= PRIORITY_PERMITS);
675 }
676
677 #[test]
678 fn an_in_flight_track_is_never_claimed_twice() {
679 let (mut q, store) = (Queue::default(), DownloadStore::new());
680 let id = qid();
681 assert_eq!(claim_priority(&mut q, &store, (1, id)), Dispatch::Spawn);
682 assert_eq!(
683 claim_priority(&mut q, &store, (1, id)),
684 Dispatch::AlreadyRunning
685 );
686 assert_eq!(q.priority_active, 1);
687 assert!(
688 q.pending.is_empty(),
689 "a duplicate request must not re-queue the track"
690 );
691 }
692
693 #[test]
694 fn playing_a_track_again_joins_the_transfer_already_running() {
695 let (mut q, store) = (Queue::default(), DownloadStore::new());
701 let (first, again) = (qid(), qid());
702 assert_eq!(claim_priority(&mut q, &store, (7, first)), Dispatch::Spawn);
703 assert_eq!(
704 claim_priority(&mut q, &store, (7, again)),
705 Dispatch::AlreadyRunning
706 );
707
708 assert_eq!(q.priority_active, 1, "one transfer, not two");
709 assert!(q.pending.is_empty());
710 assert_eq!(
711 waiters(&store, 7),
712 HashSet::from([first, again]),
713 "both entries wait on the one transfer"
714 );
715 }
716
717 #[test]
718 fn different_tracks_still_run_side_by_side() {
719 let (mut q, store) = (Queue::default(), DownloadStore::new());
721 assert_eq!(claim_priority(&mut q, &store, (1, qid())), Dispatch::Spawn);
722 assert_eq!(claim_priority(&mut q, &store, (2, qid())), Dispatch::Spawn);
723 assert_eq!(q.priority_active, 2);
724 }
725
726 #[test]
727 fn requeued_priority_item_goes_to_the_head_of_the_queue() {
728 let (mut q, store) = (Queue::default(), DownloadStore::new());
729 q.pending.push_back((9, qid()));
730 for i in 0..PRIORITY_PERMITS {
731 claim_priority(&mut q, &store, (i as i64, qid()));
732 }
733
734 let wanted = qid();
735 assert_eq!(
736 claim_priority(&mut q, &store, (7, wanted)),
737 Dispatch::Requeued
738 );
739 assert_eq!(q.pending.front().map(|(_, id)| *id), Some(wanted));
740 }
741
742 #[test]
743 fn claiming_removes_a_duplicate_queue_entry() {
744 let (mut q, store) = (Queue::default(), DownloadStore::new());
745 let id = qid();
746 q.pending.push_back((1, id));
747 q.pending.push_back((2, qid()));
748
749 assert_eq!(claim_priority(&mut q, &store, (1, id)), Dispatch::Spawn);
750 assert_eq!(
751 q.pending.len(),
752 1,
753 "the pool must not also pick up the claimed track"
754 );
755 }
756
757 #[test]
758 fn the_track_under_the_cursor_goes_first_and_goes_alone() {
759 let (a, b, c) = (qid(), qid(), qid());
760 let (mut q, store) = (Queue::default(), DownloadStore::new());
761 q.pending.extend([(1, a), (2, b), (3, c)]);
762
763 assert_eq!(next_item(&mut q, &store, Some((3, c))), Some((3, Some(c))));
765 store.claim(3, Some(c));
766
767 assert_eq!(next_item(&mut q, &store, Some((3, c))), None);
769 assert_eq!(q.pending.len(), 2, "the rest wait their turn");
770
771 let _ = downloads::withdraw(&store, 3);
773 assert_eq!(next_item(&mut q, &store, None), Some((1, Some(a))));
774 assert_eq!(next_item(&mut q, &store, None), Some((2, Some(b))));
775 }
776
777 #[test]
778 fn a_cursor_with_nothing_queued_for_it_holds_nothing_up() {
779 let (a, elsewhere) = (qid(), qid());
780 let (mut q, store) = (Queue::default(), DownloadStore::new());
781 q.pending.push_back((1, a));
782 assert_eq!(
783 next_item(&mut q, &store, Some((9, elsewhere))),
784 Some((1, Some(a)))
785 );
786 }
787
788 #[test]
789 fn a_track_wanted_only_in_the_cache_waits_behind_the_playlist() {
790 let a = qid();
791 let (mut q, store) = (Queue::default(), DownloadStore::new());
792 q.cache.push_back(5);
793 q.pending.push_back((1, a));
794 assert_eq!(next_item(&mut q, &store, None), Some((1, Some(a))));
795 assert_eq!(next_item(&mut q, &store, None), Some((5, None)));
796 }
797
798 #[test]
799 fn the_queue_lets_go_of_what_the_playlist_no_longer_holds() {
800 let (old, kept, waiter) = (qid(), qid(), qid());
803 let (mut q, store) = (Queue::default(), DownloadStore::new());
804 q.pending.extend([(1, old), (2, kept)]);
805 store.claim(3, Some(waiter));
806
807 sync_with(&mut q, &store, &[(2, kept)], None);
808
809 assert_eq!(q.pending, VecDeque::from([(2, kept)]));
810 assert!(waiters(&store, 3).is_empty());
811 assert!(
812 store.abandoned(3),
813 "a transfer nothing waits on any more is on its way to being stopped"
814 );
815 }
816
817 #[test]
818 fn the_queue_is_in_the_order_it_is_given() {
819 let (a, b, c) = (qid(), qid(), qid());
823 let (mut q, store) = (Queue::default(), DownloadStore::new());
824 q.pending.extend([(1, a), (2, b), (3, c)]);
825
826 assert!(
827 !sync_with(&mut q, &store, &[(3, c), (1, a), (2, b)], None),
828 "nothing new"
829 );
830 assert_eq!(q.pending, VecDeque::from([(3, c), (1, a), (2, b)]));
831
832 let d = qid();
833 assert!(sync_with(
834 &mut q,
835 &store,
836 &[(3, c), (4, d), (1, a), (2, b)],
837 None
838 ));
839 assert_eq!(q.pending, VecDeque::from([(3, c), (4, d), (1, a), (2, b)]));
840 }
841
842 #[test]
843 fn playing_from_the_middle_fetches_from_there_first() {
844 let (inner, ids, _) = queue_over(&[1, 2, 3, 4, 5]);
847 inner.state.set_cursor(Some(ids[2].1));
848 let wanted = inner.state.pending_downloads();
849 sync_with(
850 &mut inner.queue.lock(),
851 inner.state.downloads(),
852 &wanted,
853 None,
854 );
855
856 let order: Vec<i64> = inner.queue.lock().pending.iter().map(|(t, _)| *t).collect();
857 assert_eq!(order, vec![3, 4, 5, 1, 2]);
858 }
859
860 #[test]
861 fn a_new_entry_for_a_track_in_flight_waits_on_that_transfer() {
862 let (running, again) = (qid(), qid());
863 let (mut q, store) = (Queue::default(), DownloadStore::new());
864 store.claim(7, Some(running));
865
866 assert!(!sync_with(
867 &mut q,
868 &store,
869 &[(7, running), (7, again)],
870 None
871 ));
872
873 assert!(q.pending.is_empty());
874 assert_eq!(waiters(&store, 7), HashSet::from([running, again]));
875 }
876
877 #[test]
878 fn a_queue_replaced_with_the_same_track_keeps_its_transfer() {
879 let (old, new) = (qid(), qid());
882 let (mut q, store) = (Queue::default(), DownloadStore::new());
883 store.claim(7, Some(old));
884
885 sync_with(&mut q, &store, &[(7, new)], None);
886
887 assert_eq!(waiters(&store, 7), HashSet::from([new]));
888 assert!(!store.abandoned(7));
889 }
890
891 #[test]
892 fn entries_past_the_window_wait_until_it_reaches_them() {
893 let (played, next, beyond) = (qid(), qid(), qid());
894 let (mut q, store) = (Queue::default(), DownloadStore::new());
895 let wanted = [(2, next), (3, beyond), (1, played)];
896
897 sync_with(&mut q, &store, &wanted, Some(&HashSet::from([next])));
898 assert_eq!(q.pending, VecDeque::from([(2, next)]));
899
900 sync_with(&mut q, &store, &wanted, Some(&HashSet::from([beyond])));
902 assert_eq!(q.pending, VecDeque::from([(3, beyond)]));
903 }
904
905 #[test]
906 fn a_track_inside_the_window_keeps_its_transfer_when_the_window_is_recomputed() {
907 let (playing, next, beyond) = (qid(), qid(), qid());
908 let (mut q, store) = (Queue::default(), DownloadStore::new());
909 store.claim(2, Some(next));
910 store.claim(3, Some(beyond));
911 let wanted = [(1, playing), (2, next), (3, beyond)];
912
913 sync_with(
916 &mut q,
917 &store,
918 &wanted,
919 Some(&HashSet::from([playing, next])),
920 );
921
922 assert!(!store.abandoned(2), "inside the window: still wanted");
923 assert_eq!(waiters(&store, 2), HashSet::from([next]));
924 assert!(store.abandoned(3), "past the window: let go");
925 assert_eq!(q.pending, VecDeque::from([(1, playing)]));
926 }
927
928 fn queue_over(
931 tracks: &[i64],
932 ) -> (
933 Arc<Inner>,
934 Vec<(i64, QueueItemId)>,
935 crossbeam_channel::Receiver<PlayerCommand>,
936 ) {
937 crate::config::isolate_config_for_tests();
938 let item = |db_id: i64| PlaylistItem {
939 playlist_entry_id: None,
940 id: qid(),
941 db_id: Some(db_id),
942 path: std::path::PathBuf::from(format!("/cache/track-{db_id}.flac")),
943 title: format!("track-{db_id}"),
944 artist: "Artist".into(),
945 album_artist: "Artist".into(),
946 album: "Album".into(),
947 year: None,
948 codec: None,
949 track_number: None,
950 disc: None,
951 duration_ms: None,
952 state: ItemState::Pending,
953 pre_shuffle: None,
954 };
955 let items: Vec<_> = tracks.iter().map(|&db_id| item(db_id)).collect();
956 let ids: Vec<_> = items.iter().map(|i| (i.db_id.unwrap(), i.id)).collect();
957
958 let state = SharedPlayerState::new();
959 state.add_items(items);
960 state.set_cursor(Some(ids[0].1));
961
962 let (cmd_tx, cmd_rx) = crossbeam_channel::unbounded();
963 let inner = Arc::new(Inner {
964 queue: Mutex::new(Queue::default()),
965 has_work: Condvar::new(),
966 state,
967 cmd_tx,
968 last_evicted: Mutex::new(None),
969 spawned: std::sync::atomic::AtomicUsize::new(0),
970 });
971 (inner, ids, cmd_rx)
972 }
973
974 fn first_play() -> (Arc<Inner>, Vec<(i64, QueueItemId)>) {
977 let (inner, ids, _) = queue_over(&[1, 2, 3]);
978 inner.queue.lock().pending.extend(ids.iter().copied());
979 (inner, ids)
980 }
981
982 #[test]
983 fn a_cursor_set_before_its_tracks_were_queued_still_goes_first() {
984 let (inner, ids) = first_play();
985 promote_cursor(&inner, ids[0].1);
986
987 let left: Vec<_> = inner.queue.lock().pending.iter().copied().collect();
988 assert_eq!(
989 left,
990 vec![ids[2]],
991 "the cursor's track and the next went to the priority lane"
992 );
993 }
994
995 fn failed(inner: &Inner, id: QueueItemId) -> bool {
996 matches!(inner.state.item_state(id), Some(ItemState::Failed(_)))
997 }
998
999 #[test]
1003 fn a_watcher_started_over_a_playlist_fetches_it_unprompted() {
1004 let (inner, ids, _) = queue_over(&[1, 2, 3]);
1005 let watched = inner.clone();
1006 std::thread::spawn(move || follow_playlist(watched));
1007
1008 let deadline = std::time::Instant::now() + Duration::from_secs(5);
1009 while !ids.iter().all(|(_, id)| failed(&inner, *id)) && std::time::Instant::now() < deadline
1010 {
1011 std::thread::sleep(Duration::from_millis(10));
1012 }
1013 assert!(ids.iter().all(|(_, id)| failed(&inner, *id)));
1014 }
1015
1016 fn duplicate_that_fails() -> (
1019 Arc<Inner>,
1020 crossbeam_channel::Receiver<PlayerCommand>,
1021 QueueItemId,
1022 QueueItemId,
1023 ) {
1024 let (inner, ids, cmd_rx) = queue_over(&[1, 1]);
1025 let (first, again) = (ids[0].1, ids[1].1);
1026 inner.state.set_cursor(Some(again));
1027 let store = inner.state.downloads();
1028 store.claim(1, Some(first));
1029 store.join(1, Some(again));
1030 (inner, cmd_rx, first, again)
1031 }
1032
1033 #[test]
1034 fn a_failed_transfer_fails_every_entry_waiting_on_it() {
1035 let (inner, cmd_rx, first, again) = duplicate_that_fails();
1036
1037 run_download(&inner, 1);
1038
1039 for id in [first, again] {
1040 assert!(failed(&inner, id), "every waiter hears the one answer");
1041 }
1042 let sent: Vec<_> = cmd_rx.try_iter().collect();
1043 assert!(
1044 matches!(sent.as_slice(), [PlayerCommand::TrackFailed(id)] if *id == again),
1045 "the cursor's entry is told it failed, not that it is ready: {sent:?}"
1046 );
1047 assert!(!inner.state.downloads().in_flight(1));
1048 }
1049
1050 #[test]
1051 fn a_transfer_whose_first_entry_was_removed_still_answers_the_rest() {
1052 let (inner, _cmd_rx, first, again) = duplicate_that_fails();
1053 inner.state.remove_item(first);
1054 sync(&inner);
1055
1056 run_download(&inner, 1);
1057
1058 assert!(failed(&inner, again));
1059 }
1060
1061 #[test]
1062 fn a_transfer_withdrawn_with_an_entry_waiting_queues_it_again() {
1063 let (inner, ids, _) = queue_over(&[1]);
1065 let store = inner.state.downloads();
1066 store.claim(1, Some(ids[0].1));
1067
1068 let mut q = inner.queue.lock();
1069 for id in downloads::withdraw(store, 1) {
1070 q.pending.push_front((1, id));
1071 }
1072 assert_eq!(q.pending, VecDeque::from([ids[0]]));
1073 }
1074}