1use std::collections::{HashMap, HashSet, VecDeque};
2use std::panic::AssertUnwindSafe;
3use std::sync::{Arc, Mutex as StdMutex};
4
5use parking_lot::{Condvar, Mutex};
6
7use crate::config;
8use crate::player::commands::PlayerCommand;
9use crate::player::state::{LoadState, QueueItemId, SharedPlayerState};
10
11use crate::helpers::download_track;
12
13const PRIORITY_PERMITS: usize = 2;
17
18#[derive(Clone)]
26pub struct DownloadQueue {
27 inner: Arc<Inner>,
28}
29
30struct Inner {
31 queue: Mutex<Queue>,
32 has_work: Condvar,
33 state: Arc<SharedPlayerState>,
34 cmd_tx: crossbeam_channel::Sender<PlayerCommand>,
35 log_buf: Arc<StdMutex<Vec<String>>>,
36 last_evicted: Mutex<Option<std::time::Instant>>,
38 spawned: std::sync::atomic::AtomicUsize,
41}
42
43fn workers_allowed() -> usize {
46 config::Config::cached()
47 .remote
48 .download_workers
49 .clamp(1, 16)
50}
51
52fn ensure_workers(inner: &Arc<Inner>) {
54 use std::sync::atomic::Ordering;
55 let want = workers_allowed();
56 loop {
57 let have = inner.spawned.load(Ordering::Relaxed);
58 if have >= want {
59 return;
60 }
61 if inner
62 .spawned
63 .compare_exchange(have, have + 1, Ordering::Relaxed, Ordering::Relaxed)
64 .is_err()
65 {
66 continue;
67 }
68 let worker = inner.clone();
69 if let Err(e) = std::thread::Builder::new()
70 .name(format!("koan-dl-{have}"))
71 .spawn(move || worker_loop(worker, have))
72 {
73 log::error!("failed to spawn download worker {have}: {e}");
74 inner.spawned.fetch_sub(1, Ordering::Relaxed);
75 return;
76 }
77 }
78}
79
80fn next_item(q: &mut Queue, cursor: Option<(i64, QueueItemId)>) -> Option<(i64, QueueItemId)> {
84 match cursor {
85 Some((db_id, _)) if q.in_flight.contains_key(&db_id) => None,
86 Some((_, queue_id)) => q
87 .pending
88 .iter()
89 .position(|(_, qid)| *qid == queue_id)
90 .and_then(|ix| q.pending.remove(ix))
91 .or_else(|| q.pending.pop_front()),
92 None => q.pending.pop_front(),
93 }
94}
95
96fn retry_server_now() {
101 if let Some(client) = crate::helpers::subsonic_client(&config::Config::cached()) {
102 client.outage().retry_now();
103 }
104}
105
106fn cursor_download(state: &SharedPlayerState) -> Option<(i64, QueueItemId)> {
108 let id = state.cursor()?;
109 let item = state.get_item(id)?;
110 if !matches!(item.state, crate::player::state::ItemState::Pending) {
111 return None;
112 }
113 Some((item.db_id?, id))
114}
115
116#[derive(Default)]
119struct Queue {
120 pending: VecDeque<(i64, QueueItemId)>,
121 in_flight: HashMap<i64, HashSet<QueueItemId>>,
130 priority_active: usize,
131}
132
133#[derive(Debug, PartialEq, Eq)]
135enum Dispatch {
136 Spawn,
138 Requeued,
140 AlreadyRunning,
142}
143
144fn claim_priority(q: &mut Queue, item: (i64, QueueItemId)) -> Dispatch {
148 let (db_id, queue_id) = item;
149 q.pending.retain(|(_, qid)| *qid != queue_id);
150
151 if let Some(waiting) = q.in_flight.get_mut(&db_id) {
154 waiting.insert(queue_id);
155 return Dispatch::AlreadyRunning;
156 }
157 if q.priority_active >= PRIORITY_PERMITS {
158 q.pending.push_front(item);
159 return Dispatch::Requeued;
160 }
161 q.priority_active += 1;
162 q.in_flight.insert(db_id, HashSet::from([queue_id]));
163 Dispatch::Spawn
164}
165
166fn release_priority(q: &mut Queue, db_id: i64) {
167 q.in_flight.remove(&db_id);
168 q.priority_active = q.priority_active.saturating_sub(1);
169}
170
171struct Claim {
173 inner: Arc<Inner>,
174 db_id: i64,
175 priority: bool,
176}
177
178impl Drop for Claim {
179 fn drop(&mut self) {
180 let mut q = self.inner.queue.lock();
181 if self.priority {
182 release_priority(&mut q, self.db_id);
183 } else {
184 q.in_flight.remove(&self.db_id);
185 }
186 drop(q);
187 self.inner.has_work.notify_all();
189 }
190}
191
192impl DownloadQueue {
193 pub fn spawn(
195 cmd_tx: crossbeam_channel::Sender<PlayerCommand>,
196 state: Arc<SharedPlayerState>,
197 log_buf: Arc<StdMutex<Vec<String>>>,
198 ) -> Self {
199 let inner = Arc::new(Inner {
200 queue: Mutex::new(Queue::default()),
201 has_work: Condvar::new(),
202 state,
203 cmd_tx,
204 log_buf,
205 last_evicted: Mutex::new(None),
206 spawned: std::sync::atomic::AtomicUsize::new(0),
207 });
208 ensure_workers(&inner);
209
210 let trimmer = inner.clone();
211 let _ = std::thread::Builder::new()
212 .name("koan-dl-trim".into())
213 .spawn(move || trim_cache(&trimmer));
214
215 let watcher_inner = inner.clone();
216 if let Err(e) = std::thread::Builder::new()
217 .name("koan-dl-watch".into())
218 .spawn(move || cursor_watcher(watcher_inner))
219 {
220 log::error!("failed to spawn download cursor watcher: {}", e);
221 }
222
223 Self { inner }
224 }
225
226 pub fn enqueue(&self, items: Vec<(i64, QueueItemId)>) {
228 if items.is_empty() {
229 return;
230 }
231 ensure_workers(&self.inner);
232 retry_server_now();
233 self.inner.queue.lock().pending.extend(items);
234 self.inner.has_work.notify_all();
235 if let Some(cursor) = self.inner.state.cursor() {
236 promote_cursor(&self.inner, cursor);
237 }
238 }
239
240 pub fn prioritize(&self, db_id: i64, queue_id: QueueItemId) {
243 retry_server_now();
244 dispatch_priority(&self.inner, (db_id, queue_id));
245
246 let album_mates = self.inner.state.same_album_item_ids(queue_id);
247 if !album_mates.is_empty() {
248 let mate_set: HashSet<QueueItemId> = album_mates.into_iter().collect();
249 bump_to_front(&mut self.inner.queue.lock().pending, &mate_set);
250 self.inner.has_work.notify_all();
251 }
252 }
253}
254
255fn bump_to_front(pending: &mut VecDeque<(i64, QueueItemId)>, ids: &HashSet<QueueItemId>) {
257 let (front, rest): (VecDeque<_>, VecDeque<_>) =
258 pending.drain(..).partition(|(_, qid)| ids.contains(qid));
259 *pending = front;
260 pending.extend(rest);
261}
262
263fn dispatch_priority(inner: &Arc<Inner>, item: (i64, QueueItemId)) {
265 let dispatch = claim_priority(&mut inner.queue.lock(), item);
266 match dispatch {
267 Dispatch::AlreadyRunning => {}
268 Dispatch::Requeued => {
269 inner.has_work.notify_one();
270 }
271 Dispatch::Spawn => {
272 let spawn_inner = inner.clone();
273 let spawned = std::thread::Builder::new()
274 .name("koan-dl-prio".into())
275 .spawn(move || {
276 let _claim = Claim {
277 inner: spawn_inner.clone(),
278 db_id: item.0,
279 priority: true,
280 };
281 run_download(&spawn_inner, item);
282 });
283 if let Err(e) = spawned {
284 log::error!("failed to spawn priority download: {}", e);
285 let mut q = inner.queue.lock();
286 release_priority(&mut q, item.0);
287 q.pending.push_front(item);
288 drop(q);
289 inner.has_work.notify_one();
290 }
291 }
292 }
293}
294
295fn settle_waiters(inner: &Arc<Inner>, db_id: i64, downloaded: QueueItemId) {
301 let waiting: Vec<QueueItemId> = {
302 let q = inner.queue.lock();
303 q.in_flight
304 .get(&db_id)
305 .map(|ids| ids.iter().copied().filter(|id| *id != downloaded).collect())
306 .unwrap_or_default()
307 };
308 if waiting.is_empty() {
309 return;
310 }
311 let Some(item) = inner.state.get_item(downloaded) else {
312 return;
313 };
314 for id in waiting {
315 inner.state.update_paths(&[(id, item.path.clone())]);
316 inner.state.update_item_state(id, item.state.clone());
317 if inner.state.is_cursor(id) {
320 inner.cmd_tx.send(PlayerCommand::TrackReady(id)).ok();
321 }
322 }
323}
324
325fn run_download(inner: &Arc<Inner>, (db_id, queue_id): (i64, QueueItemId)) {
331 let cfg = config::Config::cached();
332 let Some(client) = crate::helpers::subsonic_client(&cfg) else {
333 crate::helpers::fail_track(
336 &inner.state,
337 &inner.cmd_tx,
338 queue_id,
339 crate::helpers::remote_unavailable(&cfg),
340 );
341 return;
342 };
343
344 let outcome = std::panic::catch_unwind(AssertUnwindSafe(|| {
345 download_track(
346 db_id,
347 queue_id,
348 &inner.cmd_tx,
349 &inner.log_buf,
350 &inner.state,
351 &cfg,
352 &client,
353 );
354 }));
355
356 if outcome.is_err() {
357 log::error!("download panicked for {:?}", queue_id);
358 crate::helpers::fail_track(
359 &inner.state,
360 &inner.cmd_tx,
361 queue_id,
362 "download panicked".into(),
363 );
364 }
365
366 settle_waiters(inner, db_id, queue_id);
369 trim_cache(inner);
370}
371
372const EVICT_EVERY: std::time::Duration = std::time::Duration::from_secs(60);
375
376fn trim_cache(inner: &Inner) {
381 {
382 let mut last = inner.last_evicted.lock();
383 if last.is_some_and(|t| t.elapsed() < EVICT_EVERY) {
384 return;
385 }
386 *last = Some(std::time::Instant::now());
387 }
388 let cfg = config::Config::cached();
389 if cfg.cache_limit_bytes().is_none() {
390 return;
391 }
392 let keep = inner
393 .state
394 .snapshot_playlist()
395 .0
396 .iter()
397 .filter_map(|i| i.db_id)
398 .collect();
399 match crate::db::pool::shared().get() {
400 Ok(db) => {
401 crate::helpers::evict_cache(&db, &cfg, &keep, false);
402 }
403 Err(e) => log::warn!("cache eviction: could not open the database: {e}"),
404 }
405}
406
407fn worker_loop(inner: Arc<Inner>, index: usize) {
420 loop {
421 let item = loop {
422 if let Some(client) = crate::helpers::subsonic_client(&config::Config::cached()) {
423 client.outage().hold();
424 }
425 let cursor = cursor_download(&inner.state);
428 let mut q = inner.queue.lock();
429 if index >= workers_allowed() {
430 inner.has_work.wait(&mut q);
431 continue;
432 }
433 match next_item(&mut q, cursor) {
434 Some(item) => {
435 match q.in_flight.get_mut(&item.0) {
438 Some(waiting) => {
439 waiting.insert(item.1);
440 }
441 None => {
442 q.in_flight.insert(item.0, HashSet::from([item.1]));
443 break item;
444 }
445 }
446 }
447 None => inner.has_work.wait(&mut q),
448 }
449 };
450 let _claim = Claim {
451 inner: inner.clone(),
452 db_id: item.0,
453 priority: false,
454 };
455 run_download(&inner, item);
456 }
457}
458
459fn cursor_watcher(inner: Arc<Inner>) {
472 let changed = crate::signal::engine_changed();
473 let mut seen = changed.generation();
474 let mut last_cursor: Option<QueueItemId> = None;
475 loop {
476 let current = inner.state.cursor();
477 if current != last_cursor {
478 last_cursor = current;
479 inner.has_work.notify_all();
480 if let Some(cursor_id) = current {
481 promote_cursor(&inner, cursor_id);
482 }
483 }
484 seen = changed.wait(seen);
485 }
486}
487
488fn promote_cursor(inner: &Arc<Inner>, cursor_id: QueueItemId) {
497 let is_pending = inner
498 .state
499 .item_load_state(cursor_id)
500 .is_some_and(|s| matches!(s, LoadState::Pending));
501 if !is_pending {
502 return;
503 }
504
505 let album_mate_ids: HashSet<QueueItemId> = inner
506 .state
507 .same_album_item_ids(cursor_id)
508 .into_iter()
509 .collect();
510
511 let mut priority_items = Vec::new();
512 {
513 let mut q = inner.queue.lock();
514 if let Some(pos) = q.pending.iter().position(|(_, qid)| *qid == cursor_id) {
515 priority_items.push(q.pending.remove(pos).expect("position just found"));
516
517 if !album_mate_ids.is_empty() {
518 bump_to_front(&mut q.pending, &album_mate_ids);
519 }
520
521 if let Some(next) = q.pending.pop_front() {
523 priority_items.push(next);
524 }
525 }
526 }
527
528 for item in priority_items {
529 dispatch_priority(inner, item);
530 }
531}
532
533pub fn shared(
543 cmd_tx: &crossbeam_channel::Sender<PlayerCommand>,
544 state: &Arc<SharedPlayerState>,
545 log_buf: Option<Arc<StdMutex<Vec<String>>>>,
546) -> &'static DownloadQueue {
547 static QUEUE: std::sync::OnceLock<DownloadQueue> = std::sync::OnceLock::new();
548 QUEUE.get_or_init(|| {
549 DownloadQueue::spawn(
550 cmd_tx.clone(),
551 state.clone(),
552 log_buf.unwrap_or_else(|| Arc::new(StdMutex::new(Vec::new()))),
553 )
554 })
555}
556
557#[cfg(test)]
558mod tests {
559 use super::*;
560
561 fn qid() -> QueueItemId {
562 QueueItemId::new()
563 }
564
565 #[test]
566 fn priority_lane_never_exceeds_its_permits() {
567 let mut q = Queue::default();
568
569 let mut spawned = 0;
571 for i in 0..500 {
572 if claim_priority(&mut q, (i, qid())) == Dispatch::Spawn {
573 spawned += 1;
574 }
575 assert!(
576 q.priority_active <= PRIORITY_PERMITS,
577 "priority lane over its permit count at iteration {}",
578 i
579 );
580 }
581
582 assert_eq!(spawned, PRIORITY_PERMITS, "only permitted claims may spawn");
583 assert_eq!(
584 q.pending.len(),
585 500 - PRIORITY_PERMITS,
586 "everything else must be queued, not dropped"
587 );
588 }
589
590 #[test]
591 fn released_permits_are_reusable() {
592 let mut q = Queue::default();
593 assert_eq!(claim_priority(&mut q, (1, qid())), Dispatch::Spawn);
594 assert_eq!(claim_priority(&mut q, (2, qid())), Dispatch::Spawn);
595 assert_eq!(claim_priority(&mut q, (3, qid())), Dispatch::Requeued);
596
597 release_priority(&mut q, 1);
598 assert_eq!(claim_priority(&mut q, (4, qid())), Dispatch::Spawn);
599 assert!(q.priority_active <= PRIORITY_PERMITS);
600 }
601
602 #[test]
603 fn an_in_flight_track_is_never_claimed_twice() {
604 let mut q = Queue::default();
605 let id = qid();
606 assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::Spawn);
607 assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::AlreadyRunning);
608 assert_eq!(q.priority_active, 1);
609 assert!(
610 q.pending.is_empty(),
611 "a duplicate request must not re-queue the track"
612 );
613 }
614
615 #[test]
616 fn playing_a_track_again_joins_the_transfer_already_running() {
617 let mut q = Queue::default();
623 let (first, again) = (qid(), qid());
624 assert_eq!(claim_priority(&mut q, (7, first)), Dispatch::Spawn);
625 assert_eq!(claim_priority(&mut q, (7, again)), Dispatch::AlreadyRunning);
626
627 assert_eq!(q.priority_active, 1, "one transfer, not two");
628 assert!(q.pending.is_empty());
629 assert_eq!(
630 q.in_flight.get(&7),
631 Some(&HashSet::from([first, again])),
632 "both entries wait on the one transfer"
633 );
634 }
635
636 #[test]
637 fn a_worker_picking_up_a_duplicate_waits_on_the_running_one() {
638 let mut q = Queue::default();
640 let (running, queued) = (qid(), qid());
641 assert_eq!(claim_priority(&mut q, (7, running)), Dispatch::Spawn);
642
643 match q.in_flight.get_mut(&7) {
645 Some(waiting) => {
646 waiting.insert(queued);
647 }
648 None => panic!("the track should already be claimed"),
649 }
650
651 assert_eq!(
652 q.in_flight.get(&7),
653 Some(&HashSet::from([running, queued])),
654 "the queued entry waits rather than starting a second transfer"
655 );
656 }
657
658 #[test]
659 fn different_tracks_still_run_side_by_side() {
660 let mut q = Queue::default();
662 assert_eq!(claim_priority(&mut q, (1, qid())), Dispatch::Spawn);
663 assert_eq!(claim_priority(&mut q, (2, qid())), Dispatch::Spawn);
664 assert_eq!(q.priority_active, 2);
665 }
666
667 #[test]
668 fn requeued_priority_item_goes_to_the_head_of_the_queue() {
669 let mut q = Queue::default();
670 q.pending.push_back((9, qid()));
671 for i in 0..PRIORITY_PERMITS {
672 claim_priority(&mut q, (i as i64, qid()));
673 }
674
675 let wanted = qid();
676 assert_eq!(claim_priority(&mut q, (7, wanted)), Dispatch::Requeued);
677 assert_eq!(q.pending.front().map(|(_, id)| *id), Some(wanted));
678 }
679
680 #[test]
681 fn claiming_removes_a_duplicate_queue_entry() {
682 let mut q = Queue::default();
683 let id = qid();
684 q.pending.push_back((1, id));
685 q.pending.push_back((2, qid()));
686
687 assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::Spawn);
688 assert_eq!(
689 q.pending.len(),
690 1,
691 "the pool must not also pick up the claimed track"
692 );
693 }
694
695 #[test]
696 fn bump_to_front_preserves_relative_order() {
697 let (a, b, c, d) = (qid(), qid(), qid(), qid());
698 let mut pending: VecDeque<(i64, QueueItemId)> =
699 [(1, a), (2, b), (3, c), (4, d)].into_iter().collect();
700 let mates: HashSet<QueueItemId> = [b, d].into_iter().collect();
701
702 bump_to_front(&mut pending, &mates);
703
704 let order: Vec<QueueItemId> = pending.iter().map(|(_, id)| *id).collect();
705 assert_eq!(order, vec![b, d, a, c]);
706 }
707
708 #[test]
709 fn the_track_under_the_cursor_goes_first_and_goes_alone() {
710 let (a, b, c) = (qid(), qid(), qid());
711 let mut q = Queue::default();
712 q.pending.extend([(1, a), (2, b), (3, c)]);
713
714 assert_eq!(next_item(&mut q, Some((3, c))), Some((3, c)));
716 q.in_flight.insert(3, HashSet::from([c]));
717
718 assert_eq!(next_item(&mut q, Some((3, c))), None);
720 assert_eq!(q.pending.len(), 2, "the rest wait their turn");
721
722 q.in_flight.remove(&3);
724 assert_eq!(next_item(&mut q, None), Some((1, a)));
725 assert_eq!(next_item(&mut q, None), Some((2, b)));
726 }
727
728 #[test]
729 fn a_cursor_with_nothing_queued_for_it_holds_nothing_up() {
730 let (a, elsewhere) = (qid(), qid());
731 let mut q = Queue::default();
732 q.pending.push_back((1, a));
733 assert_eq!(next_item(&mut q, Some((9, elsewhere))), Some((1, a)));
734 }
735
736 fn first_play() -> (Arc<Inner>, Vec<(i64, QueueItemId)>) {
739 crate::config::isolate_config_for_tests();
740 let item = |title: &str, db_id: i64| crate::player::state::PlaylistItem {
741 playlist_entry_id: None,
742 id: qid(),
743 db_id: Some(db_id),
744 path: std::path::PathBuf::from(format!("/cache/{title}.flac")),
745 title: title.into(),
746 artist: "Artist".into(),
747 album_artist: "Artist".into(),
748 album: "Album".into(),
749 year: None,
750 codec: None,
751 track_number: None,
752 disc: None,
753 duration_ms: None,
754 state: crate::player::state::ItemState::Pending,
755 };
756 let items = vec![item("one", 1), item("two", 2), item("three", 3)];
757 let ids: Vec<_> = items.iter().map(|i| (i.db_id.unwrap(), i.id)).collect();
758
759 let state = SharedPlayerState::new();
762 state.add_items(items);
763 state.set_cursor(Some(ids[0].1));
764
765 let (cmd_tx, _cmd_rx) = crossbeam_channel::unbounded();
766 let inner = Arc::new(Inner {
767 queue: Mutex::new(Queue::default()),
768 has_work: Condvar::new(),
769 state,
770 cmd_tx,
771 log_buf: Arc::new(StdMutex::new(Vec::new())),
772 last_evicted: Mutex::new(None),
773 spawned: std::sync::atomic::AtomicUsize::new(0),
774 });
775 inner.queue.lock().pending.extend(ids.iter().copied());
776 (inner, ids)
777 }
778
779 #[test]
780 fn a_cursor_set_before_its_tracks_were_queued_still_goes_first() {
781 let (inner, ids) = first_play();
782 promote_cursor(&inner, ids[0].1);
783
784 let left: Vec<_> = inner.queue.lock().pending.iter().copied().collect();
785 assert_eq!(
786 left,
787 vec![ids[2]],
788 "the cursor's track and the next went to the priority lane"
789 );
790 }
791
792 #[test]
795 fn a_watcher_started_after_the_cursor_moved_still_promotes_it() {
796 let (inner, ids) = first_play();
797 let watched = inner.clone();
798 std::thread::spawn(move || cursor_watcher(watched));
799
800 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
801 while inner.queue.lock().pending.len() > 1 && std::time::Instant::now() < deadline {
802 std::thread::sleep(std::time::Duration::from_millis(10));
803 }
804 let left: Vec<_> = inner.queue.lock().pending.iter().copied().collect();
805 assert_eq!(left, vec![ids[2]], "promoted without waiting for a change");
806 }
807}