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>) {
466 let changed = crate::signal::engine_changed();
467 let mut seen = changed.generation();
468 let mut last_cursor: Option<QueueItemId> = None;
469 loop {
470 seen = changed.wait(seen);
471
472 let current = inner.state.cursor();
473 if current == last_cursor {
474 continue;
475 }
476 last_cursor = current;
477 inner.has_work.notify_all();
478
479 if let Some(cursor_id) = current {
480 promote_cursor(&inner, cursor_id);
481 }
482 }
483}
484
485fn promote_cursor(inner: &Arc<Inner>, cursor_id: QueueItemId) {
494 let is_pending = inner
495 .state
496 .item_load_state(cursor_id)
497 .is_some_and(|s| matches!(s, LoadState::Pending));
498 if !is_pending {
499 return;
500 }
501
502 let album_mate_ids: HashSet<QueueItemId> = inner
503 .state
504 .same_album_item_ids(cursor_id)
505 .into_iter()
506 .collect();
507
508 let mut priority_items = Vec::new();
509 {
510 let mut q = inner.queue.lock();
511 if let Some(pos) = q.pending.iter().position(|(_, qid)| *qid == cursor_id) {
512 priority_items.push(q.pending.remove(pos).expect("position just found"));
513
514 if !album_mate_ids.is_empty() {
515 bump_to_front(&mut q.pending, &album_mate_ids);
516 }
517
518 if let Some(next) = q.pending.pop_front() {
520 priority_items.push(next);
521 }
522 }
523 }
524
525 for item in priority_items {
526 dispatch_priority(inner, item);
527 }
528}
529
530pub fn shared(
540 cmd_tx: &crossbeam_channel::Sender<PlayerCommand>,
541 state: &Arc<SharedPlayerState>,
542 log_buf: Option<Arc<StdMutex<Vec<String>>>>,
543) -> &'static DownloadQueue {
544 static QUEUE: std::sync::OnceLock<DownloadQueue> = std::sync::OnceLock::new();
545 QUEUE.get_or_init(|| {
546 DownloadQueue::spawn(
547 cmd_tx.clone(),
548 state.clone(),
549 log_buf.unwrap_or_else(|| Arc::new(StdMutex::new(Vec::new()))),
550 )
551 })
552}
553
554#[cfg(test)]
555mod tests {
556 use super::*;
557
558 fn qid() -> QueueItemId {
559 QueueItemId::new()
560 }
561
562 #[test]
563 fn priority_lane_never_exceeds_its_permits() {
564 let mut q = Queue::default();
565
566 let mut spawned = 0;
568 for i in 0..500 {
569 if claim_priority(&mut q, (i, qid())) == Dispatch::Spawn {
570 spawned += 1;
571 }
572 assert!(
573 q.priority_active <= PRIORITY_PERMITS,
574 "priority lane over its permit count at iteration {}",
575 i
576 );
577 }
578
579 assert_eq!(spawned, PRIORITY_PERMITS, "only permitted claims may spawn");
580 assert_eq!(
581 q.pending.len(),
582 500 - PRIORITY_PERMITS,
583 "everything else must be queued, not dropped"
584 );
585 }
586
587 #[test]
588 fn released_permits_are_reusable() {
589 let mut q = Queue::default();
590 assert_eq!(claim_priority(&mut q, (1, qid())), Dispatch::Spawn);
591 assert_eq!(claim_priority(&mut q, (2, qid())), Dispatch::Spawn);
592 assert_eq!(claim_priority(&mut q, (3, qid())), Dispatch::Requeued);
593
594 release_priority(&mut q, 1);
595 assert_eq!(claim_priority(&mut q, (4, qid())), Dispatch::Spawn);
596 assert!(q.priority_active <= PRIORITY_PERMITS);
597 }
598
599 #[test]
600 fn an_in_flight_track_is_never_claimed_twice() {
601 let mut q = Queue::default();
602 let id = qid();
603 assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::Spawn);
604 assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::AlreadyRunning);
605 assert_eq!(q.priority_active, 1);
606 assert!(
607 q.pending.is_empty(),
608 "a duplicate request must not re-queue the track"
609 );
610 }
611
612 #[test]
613 fn playing_a_track_again_joins_the_transfer_already_running() {
614 let mut q = Queue::default();
620 let (first, again) = (qid(), qid());
621 assert_eq!(claim_priority(&mut q, (7, first)), Dispatch::Spawn);
622 assert_eq!(claim_priority(&mut q, (7, again)), Dispatch::AlreadyRunning);
623
624 assert_eq!(q.priority_active, 1, "one transfer, not two");
625 assert!(q.pending.is_empty());
626 assert_eq!(
627 q.in_flight.get(&7),
628 Some(&HashSet::from([first, again])),
629 "both entries wait on the one transfer"
630 );
631 }
632
633 #[test]
634 fn a_worker_picking_up_a_duplicate_waits_on_the_running_one() {
635 let mut q = Queue::default();
637 let (running, queued) = (qid(), qid());
638 assert_eq!(claim_priority(&mut q, (7, running)), Dispatch::Spawn);
639
640 match q.in_flight.get_mut(&7) {
642 Some(waiting) => {
643 waiting.insert(queued);
644 }
645 None => panic!("the track should already be claimed"),
646 }
647
648 assert_eq!(
649 q.in_flight.get(&7),
650 Some(&HashSet::from([running, queued])),
651 "the queued entry waits rather than starting a second transfer"
652 );
653 }
654
655 #[test]
656 fn different_tracks_still_run_side_by_side() {
657 let mut q = Queue::default();
659 assert_eq!(claim_priority(&mut q, (1, qid())), Dispatch::Spawn);
660 assert_eq!(claim_priority(&mut q, (2, qid())), Dispatch::Spawn);
661 assert_eq!(q.priority_active, 2);
662 }
663
664 #[test]
665 fn requeued_priority_item_goes_to_the_head_of_the_queue() {
666 let mut q = Queue::default();
667 q.pending.push_back((9, qid()));
668 for i in 0..PRIORITY_PERMITS {
669 claim_priority(&mut q, (i as i64, qid()));
670 }
671
672 let wanted = qid();
673 assert_eq!(claim_priority(&mut q, (7, wanted)), Dispatch::Requeued);
674 assert_eq!(q.pending.front().map(|(_, id)| *id), Some(wanted));
675 }
676
677 #[test]
678 fn claiming_removes_a_duplicate_queue_entry() {
679 let mut q = Queue::default();
680 let id = qid();
681 q.pending.push_back((1, id));
682 q.pending.push_back((2, qid()));
683
684 assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::Spawn);
685 assert_eq!(
686 q.pending.len(),
687 1,
688 "the pool must not also pick up the claimed track"
689 );
690 }
691
692 #[test]
693 fn bump_to_front_preserves_relative_order() {
694 let (a, b, c, d) = (qid(), qid(), qid(), qid());
695 let mut pending: VecDeque<(i64, QueueItemId)> =
696 [(1, a), (2, b), (3, c), (4, d)].into_iter().collect();
697 let mates: HashSet<QueueItemId> = [b, d].into_iter().collect();
698
699 bump_to_front(&mut pending, &mates);
700
701 let order: Vec<QueueItemId> = pending.iter().map(|(_, id)| *id).collect();
702 assert_eq!(order, vec![b, d, a, c]);
703 }
704
705 #[test]
706 fn the_track_under_the_cursor_goes_first_and_goes_alone() {
707 let (a, b, c) = (qid(), qid(), qid());
708 let mut q = Queue::default();
709 q.pending.extend([(1, a), (2, b), (3, c)]);
710
711 assert_eq!(next_item(&mut q, Some((3, c))), Some((3, c)));
713 q.in_flight.insert(3, HashSet::from([c]));
714
715 assert_eq!(next_item(&mut q, Some((3, c))), None);
717 assert_eq!(q.pending.len(), 2, "the rest wait their turn");
718
719 q.in_flight.remove(&3);
721 assert_eq!(next_item(&mut q, None), Some((1, a)));
722 assert_eq!(next_item(&mut q, None), Some((2, b)));
723 }
724
725 #[test]
726 fn a_cursor_with_nothing_queued_for_it_holds_nothing_up() {
727 let (a, elsewhere) = (qid(), qid());
728 let mut q = Queue::default();
729 q.pending.push_back((1, a));
730 assert_eq!(next_item(&mut q, Some((9, elsewhere))), Some((1, a)));
731 }
732
733 #[test]
734 fn a_cursor_set_before_its_tracks_were_queued_still_goes_first() {
735 crate::config::isolate_config_for_tests();
736 let item = |title: &str, db_id: i64| crate::player::state::PlaylistItem {
737 playlist_entry_id: None,
738 id: qid(),
739 db_id: Some(db_id),
740 path: std::path::PathBuf::from(format!("/cache/{title}.flac")),
741 title: title.into(),
742 artist: "Artist".into(),
743 album_artist: "Artist".into(),
744 album: "Album".into(),
745 year: None,
746 codec: None,
747 track_number: None,
748 disc: None,
749 duration_ms: None,
750 state: crate::player::state::ItemState::Pending,
751 };
752 let items = vec![item("one", 1), item("two", 2), item("three", 3)];
753 let ids: Vec<_> = items.iter().map(|i| (i.db_id.unwrap(), i.id)).collect();
754
755 let state = SharedPlayerState::new();
758 state.add_items(items);
759 state.set_cursor(Some(ids[0].1));
760
761 let (cmd_tx, _cmd_rx) = crossbeam_channel::unbounded();
762 let inner = Arc::new(Inner {
763 queue: Mutex::new(Queue::default()),
764 has_work: Condvar::new(),
765 state,
766 cmd_tx,
767 log_buf: Arc::new(StdMutex::new(Vec::new())),
768 last_evicted: Mutex::new(None),
769 spawned: std::sync::atomic::AtomicUsize::new(0),
770 });
771 inner.queue.lock().pending.extend(ids.iter().copied());
772 promote_cursor(&inner, ids[0].1);
773
774 let left: Vec<_> = inner.queue.lock().pending.iter().copied().collect();
775 assert_eq!(
776 left,
777 vec![ids[2]],
778 "the cursor's track and the next went to the priority lane"
779 );
780 }
781}