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};
10use crate::remote::client::SubsonicClient;
11
12use crate::helpers::download_track;
13
14const PRIORITY_PERMITS: usize = 2;
18
19#[derive(Clone)]
27pub struct DownloadQueue {
28 inner: Arc<Inner>,
29}
30
31struct Inner {
32 queue: Mutex<Queue>,
33 has_work: Condvar,
34 state: Arc<SharedPlayerState>,
35 cmd_tx: crossbeam_channel::Sender<PlayerCommand>,
36 log_buf: Arc<StdMutex<Vec<String>>>,
37 cfg: config::Config,
38 client: Option<Arc<SubsonicClient>>,
40 last_evicted: Mutex<Option<std::time::Instant>>,
42}
43
44#[derive(Default)]
47struct Queue {
48 pending: VecDeque<(i64, QueueItemId)>,
49 in_flight: HashMap<i64, HashSet<QueueItemId>>,
58 priority_active: usize,
59}
60
61#[derive(Debug, PartialEq, Eq)]
63enum Dispatch {
64 Spawn,
66 Requeued,
68 AlreadyRunning,
70}
71
72fn claim_priority(q: &mut Queue, item: (i64, QueueItemId)) -> Dispatch {
76 let (db_id, queue_id) = item;
77 q.pending.retain(|(_, qid)| *qid != queue_id);
78
79 if let Some(waiting) = q.in_flight.get_mut(&db_id) {
82 waiting.insert(queue_id);
83 return Dispatch::AlreadyRunning;
84 }
85 if q.priority_active >= PRIORITY_PERMITS {
86 q.pending.push_front(item);
87 return Dispatch::Requeued;
88 }
89 q.priority_active += 1;
90 q.in_flight.insert(db_id, HashSet::from([queue_id]));
91 Dispatch::Spawn
92}
93
94fn release_priority(q: &mut Queue, db_id: i64) {
95 q.in_flight.remove(&db_id);
96 q.priority_active = q.priority_active.saturating_sub(1);
97}
98
99struct Claim {
101 inner: Arc<Inner>,
102 db_id: i64,
103 priority: bool,
104}
105
106impl Drop for Claim {
107 fn drop(&mut self) {
108 let mut q = self.inner.queue.lock();
109 if self.priority {
110 release_priority(&mut q, self.db_id);
111 } else {
112 q.in_flight.remove(&self.db_id);
113 }
114 }
115}
116
117impl DownloadQueue {
118 pub fn spawn(
120 cmd_tx: crossbeam_channel::Sender<PlayerCommand>,
121 state: Arc<SharedPlayerState>,
122 log_buf: Arc<StdMutex<Vec<String>>>,
123 ) -> Self {
124 let cfg = config::Config::load().unwrap_or_default();
125 let num_workers = cfg.remote.download_workers.max(1);
126 let client = crate::helpers::subsonic_client(&cfg);
127 if client.is_none() {
128 log::info!("remote not configured — download queue will idle");
129 }
130
131 let inner = Arc::new(Inner {
132 queue: Mutex::new(Queue::default()),
133 has_work: Condvar::new(),
134 state,
135 cmd_tx,
136 log_buf,
137 cfg,
138 client,
139 last_evicted: Mutex::new(None),
140 });
141
142 for i in 0..num_workers {
143 let inner = inner.clone();
144 if let Err(e) = std::thread::Builder::new()
145 .name(format!("koan-dl-{}", i))
146 .spawn(move || worker_loop(inner))
147 {
148 log::error!("failed to spawn download worker {}: {}", i, e);
149 }
150 }
151
152 let trimmer = inner.clone();
153 let _ = std::thread::Builder::new()
154 .name("koan-dl-trim".into())
155 .spawn(move || trim_cache(&trimmer));
156
157 let watcher_inner = inner.clone();
158 if let Err(e) = std::thread::Builder::new()
159 .name("koan-dl-watch".into())
160 .spawn(move || cursor_watcher(watcher_inner))
161 {
162 log::error!("failed to spawn download cursor watcher: {}", e);
163 }
164
165 Self { inner }
166 }
167
168 pub fn enqueue(&self, items: Vec<(i64, QueueItemId)>) {
170 if items.is_empty() {
171 return;
172 }
173 self.inner.queue.lock().pending.extend(items);
174 self.inner.has_work.notify_all();
175 }
176
177 pub fn prioritize(&self, db_id: i64, queue_id: QueueItemId) {
180 dispatch_priority(&self.inner, (db_id, queue_id));
181
182 let album_mates = self.inner.state.same_album_item_ids(queue_id);
183 if !album_mates.is_empty() {
184 let mate_set: HashSet<QueueItemId> = album_mates.into_iter().collect();
185 bump_to_front(&mut self.inner.queue.lock().pending, &mate_set);
186 self.inner.has_work.notify_all();
187 }
188 }
189}
190
191fn bump_to_front(pending: &mut VecDeque<(i64, QueueItemId)>, ids: &HashSet<QueueItemId>) {
193 let (front, rest): (VecDeque<_>, VecDeque<_>) =
194 pending.drain(..).partition(|(_, qid)| ids.contains(qid));
195 *pending = front;
196 pending.extend(rest);
197}
198
199fn dispatch_priority(inner: &Arc<Inner>, item: (i64, QueueItemId)) {
201 let dispatch = claim_priority(&mut inner.queue.lock(), item);
202 match dispatch {
203 Dispatch::AlreadyRunning => {}
204 Dispatch::Requeued => {
205 inner.has_work.notify_one();
206 }
207 Dispatch::Spawn => {
208 let spawn_inner = inner.clone();
209 let spawned = std::thread::Builder::new()
210 .name("koan-dl-prio".into())
211 .spawn(move || {
212 let _claim = Claim {
213 inner: spawn_inner.clone(),
214 db_id: item.0,
215 priority: true,
216 };
217 run_download(&spawn_inner, item);
218 });
219 if let Err(e) = spawned {
220 log::error!("failed to spawn priority download: {}", e);
221 let mut q = inner.queue.lock();
222 release_priority(&mut q, item.0);
223 q.pending.push_front(item);
224 drop(q);
225 inner.has_work.notify_one();
226 }
227 }
228 }
229}
230
231fn settle_waiters(inner: &Arc<Inner>, db_id: i64, downloaded: QueueItemId) {
237 let waiting: Vec<QueueItemId> = {
238 let q = inner.queue.lock();
239 q.in_flight
240 .get(&db_id)
241 .map(|ids| ids.iter().copied().filter(|id| *id != downloaded).collect())
242 .unwrap_or_default()
243 };
244 if waiting.is_empty() {
245 return;
246 }
247 let Some(item) = inner.state.get_item(downloaded) else {
248 return;
249 };
250 for id in waiting {
251 inner.state.update_paths(&[(id, item.path.clone())]);
252 inner.state.update_item_state(id, item.state.clone());
253 if inner.state.is_cursor(id) {
256 inner.cmd_tx.send(PlayerCommand::TrackReady(id)).ok();
257 }
258 }
259}
260
261fn run_download(inner: &Arc<Inner>, (db_id, queue_id): (i64, QueueItemId)) {
263 let Some(client) = inner.client.as_ref() else {
264 crate::helpers::fail_track(
267 &inner.state,
268 &inner.cmd_tx,
269 queue_id,
270 crate::helpers::remote_unavailable(&inner.cfg),
271 );
272 return;
273 };
274
275 let outcome = std::panic::catch_unwind(AssertUnwindSafe(|| {
276 download_track(
277 db_id,
278 queue_id,
279 &inner.cmd_tx,
280 &inner.log_buf,
281 &inner.state,
282 &inner.cfg,
283 client,
284 );
285 }));
286
287 if outcome.is_err() {
288 log::error!("download panicked for {:?}", queue_id);
289 crate::helpers::fail_track(
290 &inner.state,
291 &inner.cmd_tx,
292 queue_id,
293 "download panicked".into(),
294 );
295 }
296
297 settle_waiters(inner, db_id, queue_id);
300 trim_cache(inner);
301}
302
303const EVICT_EVERY: std::time::Duration = std::time::Duration::from_secs(60);
306
307fn trim_cache(inner: &Inner) {
312 {
313 let mut last = inner.last_evicted.lock();
314 if last.is_some_and(|t| t.elapsed() < EVICT_EVERY) {
315 return;
316 }
317 *last = Some(std::time::Instant::now());
318 }
319 let cfg = config::Config::load().unwrap_or_else(|_| inner.cfg.clone());
320 if cfg.cache_limit_bytes().is_none() {
321 return;
322 }
323 let keep = inner
324 .state
325 .snapshot_playlist()
326 .0
327 .iter()
328 .filter_map(|i| i.db_id)
329 .collect();
330 match crate::db::connection::Database::open_default() {
331 Ok(db) => {
332 crate::helpers::evict_cache(&db, &cfg, &keep, false);
333 }
334 Err(e) => log::warn!("cache eviction: could not open the database: {e}"),
335 }
336}
337
338fn worker_loop(inner: Arc<Inner>) {
344 loop {
345 let item = loop {
346 if let Some(client) = &inner.client {
347 client.outage().hold();
348 }
349 let mut q = inner.queue.lock();
350 match q.pending.pop_front() {
351 Some(item) => {
352 match q.in_flight.get_mut(&item.0) {
355 Some(waiting) => {
356 waiting.insert(item.1);
357 }
358 None => {
359 q.in_flight.insert(item.0, HashSet::from([item.1]));
360 break item;
361 }
362 }
363 }
364 None => inner.has_work.wait(&mut q),
365 }
366 };
367 let _claim = Claim {
368 inner: inner.clone(),
369 db_id: item.0,
370 priority: false,
371 };
372 run_download(&inner, item);
373 }
374}
375
376fn cursor_watcher(inner: Arc<Inner>) {
379 let changed = crate::signal::engine_changed();
380 let mut seen = changed.generation();
381 let mut last_cursor: Option<QueueItemId> = None;
382 loop {
383 seen = changed.wait(seen);
384
385 let current = inner.state.cursor();
386 if current == last_cursor {
387 continue;
388 }
389 last_cursor = current;
390
391 let Some(cursor_id) = current else {
392 continue;
393 };
394
395 let is_pending = inner
396 .state
397 .item_load_state(cursor_id)
398 .is_some_and(|s| matches!(s, LoadState::Pending));
399 if !is_pending {
400 continue;
401 }
402
403 let album_mate_ids: HashSet<QueueItemId> = inner
404 .state
405 .same_album_item_ids(cursor_id)
406 .into_iter()
407 .collect();
408
409 let mut priority_items = Vec::new();
410 {
411 let mut q = inner.queue.lock();
412 if let Some(pos) = q.pending.iter().position(|(_, qid)| *qid == cursor_id) {
413 priority_items.push(q.pending.remove(pos).expect("position just found"));
414
415 if !album_mate_ids.is_empty() {
416 bump_to_front(&mut q.pending, &album_mate_ids);
417 }
418
419 if let Some(next) = q.pending.pop_front() {
421 priority_items.push(next);
422 }
423 }
424 }
425
426 for item in priority_items {
427 dispatch_priority(&inner, item);
428 }
429 }
430}
431
432pub fn shared(
442 cmd_tx: &crossbeam_channel::Sender<PlayerCommand>,
443 state: &Arc<SharedPlayerState>,
444 log_buf: Option<Arc<StdMutex<Vec<String>>>>,
445) -> &'static DownloadQueue {
446 static QUEUE: std::sync::OnceLock<DownloadQueue> = std::sync::OnceLock::new();
447 QUEUE.get_or_init(|| {
448 DownloadQueue::spawn(
449 cmd_tx.clone(),
450 state.clone(),
451 log_buf.unwrap_or_else(|| Arc::new(StdMutex::new(Vec::new()))),
452 )
453 })
454}
455
456#[cfg(test)]
457mod tests {
458 use super::*;
459
460 fn qid() -> QueueItemId {
461 QueueItemId::new()
462 }
463
464 #[test]
465 fn priority_lane_never_exceeds_its_permits() {
466 let mut q = Queue::default();
467
468 let mut spawned = 0;
470 for i in 0..500 {
471 if claim_priority(&mut q, (i, qid())) == Dispatch::Spawn {
472 spawned += 1;
473 }
474 assert!(
475 q.priority_active <= PRIORITY_PERMITS,
476 "priority lane over its permit count at iteration {}",
477 i
478 );
479 }
480
481 assert_eq!(spawned, PRIORITY_PERMITS, "only permitted claims may spawn");
482 assert_eq!(
483 q.pending.len(),
484 500 - PRIORITY_PERMITS,
485 "everything else must be queued, not dropped"
486 );
487 }
488
489 #[test]
490 fn released_permits_are_reusable() {
491 let mut q = Queue::default();
492 assert_eq!(claim_priority(&mut q, (1, qid())), Dispatch::Spawn);
493 assert_eq!(claim_priority(&mut q, (2, qid())), Dispatch::Spawn);
494 assert_eq!(claim_priority(&mut q, (3, qid())), Dispatch::Requeued);
495
496 release_priority(&mut q, 1);
497 assert_eq!(claim_priority(&mut q, (4, qid())), Dispatch::Spawn);
498 assert!(q.priority_active <= PRIORITY_PERMITS);
499 }
500
501 #[test]
502 fn an_in_flight_track_is_never_claimed_twice() {
503 let mut q = Queue::default();
504 let id = qid();
505 assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::Spawn);
506 assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::AlreadyRunning);
507 assert_eq!(q.priority_active, 1);
508 assert!(
509 q.pending.is_empty(),
510 "a duplicate request must not re-queue the track"
511 );
512 }
513
514 #[test]
515 fn playing_a_track_again_joins_the_transfer_already_running() {
516 let mut q = Queue::default();
522 let (first, again) = (qid(), qid());
523 assert_eq!(claim_priority(&mut q, (7, first)), Dispatch::Spawn);
524 assert_eq!(claim_priority(&mut q, (7, again)), Dispatch::AlreadyRunning);
525
526 assert_eq!(q.priority_active, 1, "one transfer, not two");
527 assert!(q.pending.is_empty());
528 assert_eq!(
529 q.in_flight.get(&7),
530 Some(&HashSet::from([first, again])),
531 "both entries wait on the one transfer"
532 );
533 }
534
535 #[test]
536 fn a_worker_picking_up_a_duplicate_waits_on_the_running_one() {
537 let mut q = Queue::default();
539 let (running, queued) = (qid(), qid());
540 assert_eq!(claim_priority(&mut q, (7, running)), Dispatch::Spawn);
541
542 match q.in_flight.get_mut(&7) {
544 Some(waiting) => {
545 waiting.insert(queued);
546 }
547 None => panic!("the track should already be claimed"),
548 }
549
550 assert_eq!(
551 q.in_flight.get(&7),
552 Some(&HashSet::from([running, queued])),
553 "the queued entry waits rather than starting a second transfer"
554 );
555 }
556
557 #[test]
558 fn different_tracks_still_run_side_by_side() {
559 let mut q = Queue::default();
561 assert_eq!(claim_priority(&mut q, (1, qid())), Dispatch::Spawn);
562 assert_eq!(claim_priority(&mut q, (2, qid())), Dispatch::Spawn);
563 assert_eq!(q.priority_active, 2);
564 }
565
566 #[test]
567 fn requeued_priority_item_goes_to_the_head_of_the_queue() {
568 let mut q = Queue::default();
569 q.pending.push_back((9, qid()));
570 for i in 0..PRIORITY_PERMITS {
571 claim_priority(&mut q, (i as i64, qid()));
572 }
573
574 let wanted = qid();
575 assert_eq!(claim_priority(&mut q, (7, wanted)), Dispatch::Requeued);
576 assert_eq!(q.pending.front().map(|(_, id)| *id), Some(wanted));
577 }
578
579 #[test]
580 fn claiming_removes_a_duplicate_queue_entry() {
581 let mut q = Queue::default();
582 let id = qid();
583 q.pending.push_back((1, id));
584 q.pending.push_back((2, qid()));
585
586 assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::Spawn);
587 assert_eq!(
588 q.pending.len(),
589 1,
590 "the pool must not also pick up the claimed track"
591 );
592 }
593
594 #[test]
595 fn bump_to_front_preserves_relative_order() {
596 let (a, b, c, d) = (qid(), qid(), qid(), qid());
597 let mut pending: VecDeque<(i64, QueueItemId)> =
598 [(1, a), (2, b), (3, c), (4, d)].into_iter().collect();
599 let mates: HashSet<QueueItemId> = [b, d].into_iter().collect();
600
601 bump_to_front(&mut pending, &mates);
602
603 let order: Vec<QueueItemId> = pending.iter().map(|(_, id)| *id).collect();
604 assert_eq!(order, vec![b, d, a, c]);
605 }
606}