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 cursor_download(state: &SharedPlayerState) -> Option<(i64, QueueItemId)> {
98 let id = state.cursor()?;
99 let item = state.get_item(id)?;
100 if !matches!(item.state, crate::player::state::ItemState::Pending) {
101 return None;
102 }
103 Some((item.db_id?, id))
104}
105
106#[derive(Default)]
109struct Queue {
110 pending: VecDeque<(i64, QueueItemId)>,
111 in_flight: HashMap<i64, HashSet<QueueItemId>>,
120 priority_active: usize,
121}
122
123#[derive(Debug, PartialEq, Eq)]
125enum Dispatch {
126 Spawn,
128 Requeued,
130 AlreadyRunning,
132}
133
134fn claim_priority(q: &mut Queue, item: (i64, QueueItemId)) -> Dispatch {
138 let (db_id, queue_id) = item;
139 q.pending.retain(|(_, qid)| *qid != queue_id);
140
141 if let Some(waiting) = q.in_flight.get_mut(&db_id) {
144 waiting.insert(queue_id);
145 return Dispatch::AlreadyRunning;
146 }
147 if q.priority_active >= PRIORITY_PERMITS {
148 q.pending.push_front(item);
149 return Dispatch::Requeued;
150 }
151 q.priority_active += 1;
152 q.in_flight.insert(db_id, HashSet::from([queue_id]));
153 Dispatch::Spawn
154}
155
156fn release_priority(q: &mut Queue, db_id: i64) {
157 q.in_flight.remove(&db_id);
158 q.priority_active = q.priority_active.saturating_sub(1);
159}
160
161struct Claim {
163 inner: Arc<Inner>,
164 db_id: i64,
165 priority: bool,
166}
167
168impl Drop for Claim {
169 fn drop(&mut self) {
170 let mut q = self.inner.queue.lock();
171 if self.priority {
172 release_priority(&mut q, self.db_id);
173 } else {
174 q.in_flight.remove(&self.db_id);
175 }
176 drop(q);
177 self.inner.has_work.notify_all();
179 }
180}
181
182impl DownloadQueue {
183 pub fn spawn(
185 cmd_tx: crossbeam_channel::Sender<PlayerCommand>,
186 state: Arc<SharedPlayerState>,
187 log_buf: Arc<StdMutex<Vec<String>>>,
188 ) -> Self {
189 let inner = Arc::new(Inner {
190 queue: Mutex::new(Queue::default()),
191 has_work: Condvar::new(),
192 state,
193 cmd_tx,
194 log_buf,
195 last_evicted: Mutex::new(None),
196 spawned: std::sync::atomic::AtomicUsize::new(0),
197 });
198 ensure_workers(&inner);
199
200 let trimmer = inner.clone();
201 let _ = std::thread::Builder::new()
202 .name("koan-dl-trim".into())
203 .spawn(move || trim_cache(&trimmer));
204
205 let watcher_inner = inner.clone();
206 if let Err(e) = std::thread::Builder::new()
207 .name("koan-dl-watch".into())
208 .spawn(move || cursor_watcher(watcher_inner))
209 {
210 log::error!("failed to spawn download cursor watcher: {}", e);
211 }
212
213 Self { inner }
214 }
215
216 pub fn enqueue(&self, items: Vec<(i64, QueueItemId)>) {
218 if items.is_empty() {
219 return;
220 }
221 ensure_workers(&self.inner);
222 self.inner.queue.lock().pending.extend(items);
223 self.inner.has_work.notify_all();
224 }
225
226 pub fn prioritize(&self, db_id: i64, queue_id: QueueItemId) {
229 dispatch_priority(&self.inner, (db_id, queue_id));
230
231 let album_mates = self.inner.state.same_album_item_ids(queue_id);
232 if !album_mates.is_empty() {
233 let mate_set: HashSet<QueueItemId> = album_mates.into_iter().collect();
234 bump_to_front(&mut self.inner.queue.lock().pending, &mate_set);
235 self.inner.has_work.notify_all();
236 }
237 }
238}
239
240fn bump_to_front(pending: &mut VecDeque<(i64, QueueItemId)>, ids: &HashSet<QueueItemId>) {
242 let (front, rest): (VecDeque<_>, VecDeque<_>) =
243 pending.drain(..).partition(|(_, qid)| ids.contains(qid));
244 *pending = front;
245 pending.extend(rest);
246}
247
248fn dispatch_priority(inner: &Arc<Inner>, item: (i64, QueueItemId)) {
250 let dispatch = claim_priority(&mut inner.queue.lock(), item);
251 match dispatch {
252 Dispatch::AlreadyRunning => {}
253 Dispatch::Requeued => {
254 inner.has_work.notify_one();
255 }
256 Dispatch::Spawn => {
257 let spawn_inner = inner.clone();
258 let spawned = std::thread::Builder::new()
259 .name("koan-dl-prio".into())
260 .spawn(move || {
261 let _claim = Claim {
262 inner: spawn_inner.clone(),
263 db_id: item.0,
264 priority: true,
265 };
266 run_download(&spawn_inner, item);
267 });
268 if let Err(e) = spawned {
269 log::error!("failed to spawn priority download: {}", e);
270 let mut q = inner.queue.lock();
271 release_priority(&mut q, item.0);
272 q.pending.push_front(item);
273 drop(q);
274 inner.has_work.notify_one();
275 }
276 }
277 }
278}
279
280fn settle_waiters(inner: &Arc<Inner>, db_id: i64, downloaded: QueueItemId) {
286 let waiting: Vec<QueueItemId> = {
287 let q = inner.queue.lock();
288 q.in_flight
289 .get(&db_id)
290 .map(|ids| ids.iter().copied().filter(|id| *id != downloaded).collect())
291 .unwrap_or_default()
292 };
293 if waiting.is_empty() {
294 return;
295 }
296 let Some(item) = inner.state.get_item(downloaded) else {
297 return;
298 };
299 for id in waiting {
300 inner.state.update_paths(&[(id, item.path.clone())]);
301 inner.state.update_item_state(id, item.state.clone());
302 if inner.state.is_cursor(id) {
305 inner.cmd_tx.send(PlayerCommand::TrackReady(id)).ok();
306 }
307 }
308}
309
310fn run_download(inner: &Arc<Inner>, (db_id, queue_id): (i64, QueueItemId)) {
316 let cfg = config::Config::cached();
317 let Some(client) = crate::helpers::subsonic_client(&cfg) else {
318 crate::helpers::fail_track(
321 &inner.state,
322 &inner.cmd_tx,
323 queue_id,
324 crate::helpers::remote_unavailable(&cfg),
325 );
326 return;
327 };
328
329 let outcome = std::panic::catch_unwind(AssertUnwindSafe(|| {
330 download_track(
331 db_id,
332 queue_id,
333 &inner.cmd_tx,
334 &inner.log_buf,
335 &inner.state,
336 &cfg,
337 &client,
338 );
339 }));
340
341 if outcome.is_err() {
342 log::error!("download panicked for {:?}", queue_id);
343 crate::helpers::fail_track(
344 &inner.state,
345 &inner.cmd_tx,
346 queue_id,
347 "download panicked".into(),
348 );
349 }
350
351 settle_waiters(inner, db_id, queue_id);
354 trim_cache(inner);
355}
356
357const EVICT_EVERY: std::time::Duration = std::time::Duration::from_secs(60);
360
361fn trim_cache(inner: &Inner) {
366 {
367 let mut last = inner.last_evicted.lock();
368 if last.is_some_and(|t| t.elapsed() < EVICT_EVERY) {
369 return;
370 }
371 *last = Some(std::time::Instant::now());
372 }
373 let cfg = config::Config::cached();
374 if cfg.cache_limit_bytes().is_none() {
375 return;
376 }
377 let keep = inner
378 .state
379 .snapshot_playlist()
380 .0
381 .iter()
382 .filter_map(|i| i.db_id)
383 .collect();
384 match crate::db::connection::Database::open_default() {
385 Ok(db) => {
386 crate::helpers::evict_cache(&db, &cfg, &keep, false);
387 }
388 Err(e) => log::warn!("cache eviction: could not open the database: {e}"),
389 }
390}
391
392fn worker_loop(inner: Arc<Inner>, index: usize) {
405 loop {
406 let item = loop {
407 if let Some(client) = crate::helpers::subsonic_client(&config::Config::cached()) {
408 client.outage().hold();
409 }
410 let cursor = cursor_download(&inner.state);
413 let mut q = inner.queue.lock();
414 if index >= workers_allowed() {
415 inner.has_work.wait(&mut q);
416 continue;
417 }
418 match next_item(&mut q, cursor) {
419 Some(item) => {
420 match q.in_flight.get_mut(&item.0) {
423 Some(waiting) => {
424 waiting.insert(item.1);
425 }
426 None => {
427 q.in_flight.insert(item.0, HashSet::from([item.1]));
428 break item;
429 }
430 }
431 }
432 None => inner.has_work.wait(&mut q),
433 }
434 };
435 let _claim = Claim {
436 inner: inner.clone(),
437 db_id: item.0,
438 priority: false,
439 };
440 run_download(&inner, item);
441 }
442}
443
444fn cursor_watcher(inner: Arc<Inner>) {
447 let changed = crate::signal::engine_changed();
448 let mut seen = changed.generation();
449 let mut last_cursor: Option<QueueItemId> = None;
450 loop {
451 seen = changed.wait(seen);
452
453 let current = inner.state.cursor();
454 if current == last_cursor {
455 continue;
456 }
457 last_cursor = current;
458
459 let Some(cursor_id) = current else {
460 continue;
461 };
462
463 let is_pending = inner
464 .state
465 .item_load_state(cursor_id)
466 .is_some_and(|s| matches!(s, LoadState::Pending));
467 if !is_pending {
468 continue;
469 }
470
471 let album_mate_ids: HashSet<QueueItemId> = inner
472 .state
473 .same_album_item_ids(cursor_id)
474 .into_iter()
475 .collect();
476
477 let mut priority_items = Vec::new();
478 {
479 let mut q = inner.queue.lock();
480 if let Some(pos) = q.pending.iter().position(|(_, qid)| *qid == cursor_id) {
481 priority_items.push(q.pending.remove(pos).expect("position just found"));
482
483 if !album_mate_ids.is_empty() {
484 bump_to_front(&mut q.pending, &album_mate_ids);
485 }
486
487 if let Some(next) = q.pending.pop_front() {
489 priority_items.push(next);
490 }
491 }
492 }
493
494 for item in priority_items {
495 dispatch_priority(&inner, item);
496 }
497 }
498}
499
500pub fn shared(
510 cmd_tx: &crossbeam_channel::Sender<PlayerCommand>,
511 state: &Arc<SharedPlayerState>,
512 log_buf: Option<Arc<StdMutex<Vec<String>>>>,
513) -> &'static DownloadQueue {
514 static QUEUE: std::sync::OnceLock<DownloadQueue> = std::sync::OnceLock::new();
515 QUEUE.get_or_init(|| {
516 DownloadQueue::spawn(
517 cmd_tx.clone(),
518 state.clone(),
519 log_buf.unwrap_or_else(|| Arc::new(StdMutex::new(Vec::new()))),
520 )
521 })
522}
523
524#[cfg(test)]
525mod tests {
526 use super::*;
527
528 fn qid() -> QueueItemId {
529 QueueItemId::new()
530 }
531
532 #[test]
533 fn priority_lane_never_exceeds_its_permits() {
534 let mut q = Queue::default();
535
536 let mut spawned = 0;
538 for i in 0..500 {
539 if claim_priority(&mut q, (i, qid())) == Dispatch::Spawn {
540 spawned += 1;
541 }
542 assert!(
543 q.priority_active <= PRIORITY_PERMITS,
544 "priority lane over its permit count at iteration {}",
545 i
546 );
547 }
548
549 assert_eq!(spawned, PRIORITY_PERMITS, "only permitted claims may spawn");
550 assert_eq!(
551 q.pending.len(),
552 500 - PRIORITY_PERMITS,
553 "everything else must be queued, not dropped"
554 );
555 }
556
557 #[test]
558 fn released_permits_are_reusable() {
559 let mut q = Queue::default();
560 assert_eq!(claim_priority(&mut q, (1, qid())), Dispatch::Spawn);
561 assert_eq!(claim_priority(&mut q, (2, qid())), Dispatch::Spawn);
562 assert_eq!(claim_priority(&mut q, (3, qid())), Dispatch::Requeued);
563
564 release_priority(&mut q, 1);
565 assert_eq!(claim_priority(&mut q, (4, qid())), Dispatch::Spawn);
566 assert!(q.priority_active <= PRIORITY_PERMITS);
567 }
568
569 #[test]
570 fn an_in_flight_track_is_never_claimed_twice() {
571 let mut q = Queue::default();
572 let id = qid();
573 assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::Spawn);
574 assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::AlreadyRunning);
575 assert_eq!(q.priority_active, 1);
576 assert!(
577 q.pending.is_empty(),
578 "a duplicate request must not re-queue the track"
579 );
580 }
581
582 #[test]
583 fn playing_a_track_again_joins_the_transfer_already_running() {
584 let mut q = Queue::default();
590 let (first, again) = (qid(), qid());
591 assert_eq!(claim_priority(&mut q, (7, first)), Dispatch::Spawn);
592 assert_eq!(claim_priority(&mut q, (7, again)), Dispatch::AlreadyRunning);
593
594 assert_eq!(q.priority_active, 1, "one transfer, not two");
595 assert!(q.pending.is_empty());
596 assert_eq!(
597 q.in_flight.get(&7),
598 Some(&HashSet::from([first, again])),
599 "both entries wait on the one transfer"
600 );
601 }
602
603 #[test]
604 fn a_worker_picking_up_a_duplicate_waits_on_the_running_one() {
605 let mut q = Queue::default();
607 let (running, queued) = (qid(), qid());
608 assert_eq!(claim_priority(&mut q, (7, running)), Dispatch::Spawn);
609
610 match q.in_flight.get_mut(&7) {
612 Some(waiting) => {
613 waiting.insert(queued);
614 }
615 None => panic!("the track should already be claimed"),
616 }
617
618 assert_eq!(
619 q.in_flight.get(&7),
620 Some(&HashSet::from([running, queued])),
621 "the queued entry waits rather than starting a second transfer"
622 );
623 }
624
625 #[test]
626 fn different_tracks_still_run_side_by_side() {
627 let mut q = Queue::default();
629 assert_eq!(claim_priority(&mut q, (1, qid())), Dispatch::Spawn);
630 assert_eq!(claim_priority(&mut q, (2, qid())), Dispatch::Spawn);
631 assert_eq!(q.priority_active, 2);
632 }
633
634 #[test]
635 fn requeued_priority_item_goes_to_the_head_of_the_queue() {
636 let mut q = Queue::default();
637 q.pending.push_back((9, qid()));
638 for i in 0..PRIORITY_PERMITS {
639 claim_priority(&mut q, (i as i64, qid()));
640 }
641
642 let wanted = qid();
643 assert_eq!(claim_priority(&mut q, (7, wanted)), Dispatch::Requeued);
644 assert_eq!(q.pending.front().map(|(_, id)| *id), Some(wanted));
645 }
646
647 #[test]
648 fn claiming_removes_a_duplicate_queue_entry() {
649 let mut q = Queue::default();
650 let id = qid();
651 q.pending.push_back((1, id));
652 q.pending.push_back((2, qid()));
653
654 assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::Spawn);
655 assert_eq!(
656 q.pending.len(),
657 1,
658 "the pool must not also pick up the claimed track"
659 );
660 }
661
662 #[test]
663 fn bump_to_front_preserves_relative_order() {
664 let (a, b, c, d) = (qid(), qid(), qid(), qid());
665 let mut pending: VecDeque<(i64, QueueItemId)> =
666 [(1, a), (2, b), (3, c), (4, d)].into_iter().collect();
667 let mates: HashSet<QueueItemId> = [b, d].into_iter().collect();
668
669 bump_to_front(&mut pending, &mates);
670
671 let order: Vec<QueueItemId> = pending.iter().map(|(_, id)| *id).collect();
672 assert_eq!(order, vec![b, d, a, c]);
673 }
674
675 #[test]
676 fn the_track_under_the_cursor_goes_first_and_goes_alone() {
677 let (a, b, c) = (qid(), qid(), qid());
678 let mut q = Queue::default();
679 q.pending.extend([(1, a), (2, b), (3, c)]);
680
681 assert_eq!(next_item(&mut q, Some((3, c))), Some((3, c)));
683 q.in_flight.insert(3, HashSet::from([c]));
684
685 assert_eq!(next_item(&mut q, Some((3, c))), None);
687 assert_eq!(q.pending.len(), 2, "the rest wait their turn");
688
689 q.in_flight.remove(&3);
691 assert_eq!(next_item(&mut q, None), Some((1, a)));
692 assert_eq!(next_item(&mut q, None), Some((2, b)));
693 }
694
695 #[test]
696 fn a_cursor_with_nothing_queued_for_it_holds_nothing_up() {
697 let (a, elsewhere) = (qid(), qid());
698 let mut q = Queue::default();
699 q.pending.push_back((1, a));
700 assert_eq!(next_item(&mut q, Some((9, elsewhere))), Some((1, a)));
701 }
702}