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}
41
42#[derive(Default)]
45struct Queue {
46 pending: VecDeque<(i64, QueueItemId)>,
47 in_flight: HashMap<i64, HashSet<QueueItemId>>,
56 priority_active: usize,
57}
58
59#[derive(Debug, PartialEq, Eq)]
61enum Dispatch {
62 Spawn,
64 Requeued,
66 AlreadyRunning,
68}
69
70fn claim_priority(q: &mut Queue, item: (i64, QueueItemId)) -> Dispatch {
74 let (db_id, queue_id) = item;
75 q.pending.retain(|(_, qid)| *qid != queue_id);
76
77 if let Some(waiting) = q.in_flight.get_mut(&db_id) {
80 waiting.insert(queue_id);
81 return Dispatch::AlreadyRunning;
82 }
83 if q.priority_active >= PRIORITY_PERMITS {
84 q.pending.push_front(item);
85 return Dispatch::Requeued;
86 }
87 q.priority_active += 1;
88 q.in_flight.insert(db_id, HashSet::from([queue_id]));
89 Dispatch::Spawn
90}
91
92fn release_priority(q: &mut Queue, db_id: i64) {
93 q.in_flight.remove(&db_id);
94 q.priority_active = q.priority_active.saturating_sub(1);
95}
96
97struct Claim {
99 inner: Arc<Inner>,
100 db_id: i64,
101 priority: bool,
102}
103
104impl Drop for Claim {
105 fn drop(&mut self) {
106 let mut q = self.inner.queue.lock();
107 if self.priority {
108 release_priority(&mut q, self.db_id);
109 } else {
110 q.in_flight.remove(&self.db_id);
111 }
112 }
113}
114
115impl DownloadQueue {
116 pub fn spawn(
118 cmd_tx: crossbeam_channel::Sender<PlayerCommand>,
119 state: Arc<SharedPlayerState>,
120 log_buf: Arc<StdMutex<Vec<String>>>,
121 ) -> Self {
122 let cfg = config::Config::load().unwrap_or_default();
123 let num_workers = cfg.remote.download_workers.max(1);
124 let client = crate::helpers::subsonic_client(&cfg);
125 if client.is_none() {
126 log::info!("remote not configured — download queue will idle");
127 }
128
129 let inner = Arc::new(Inner {
130 queue: Mutex::new(Queue::default()),
131 has_work: Condvar::new(),
132 state,
133 cmd_tx,
134 log_buf,
135 cfg,
136 client,
137 });
138
139 for i in 0..num_workers {
140 let inner = inner.clone();
141 if let Err(e) = std::thread::Builder::new()
142 .name(format!("koan-dl-{}", i))
143 .spawn(move || worker_loop(inner))
144 {
145 log::error!("failed to spawn download worker {}: {}", i, e);
146 }
147 }
148
149 let watcher_inner = inner.clone();
150 if let Err(e) = std::thread::Builder::new()
151 .name("koan-dl-watch".into())
152 .spawn(move || cursor_watcher(watcher_inner))
153 {
154 log::error!("failed to spawn download cursor watcher: {}", e);
155 }
156
157 Self { inner }
158 }
159
160 pub fn enqueue(&self, items: Vec<(i64, QueueItemId)>) {
162 if items.is_empty() {
163 return;
164 }
165 self.inner.queue.lock().pending.extend(items);
166 self.inner.has_work.notify_all();
167 }
168
169 pub fn prioritize(&self, db_id: i64, queue_id: QueueItemId) {
172 dispatch_priority(&self.inner, (db_id, queue_id));
173
174 let album_mates = self.inner.state.same_album_item_ids(queue_id);
175 if !album_mates.is_empty() {
176 let mate_set: HashSet<QueueItemId> = album_mates.into_iter().collect();
177 bump_to_front(&mut self.inner.queue.lock().pending, &mate_set);
178 self.inner.has_work.notify_all();
179 }
180 }
181}
182
183fn bump_to_front(pending: &mut VecDeque<(i64, QueueItemId)>, ids: &HashSet<QueueItemId>) {
185 let (front, rest): (VecDeque<_>, VecDeque<_>) =
186 pending.drain(..).partition(|(_, qid)| ids.contains(qid));
187 *pending = front;
188 pending.extend(rest);
189}
190
191fn dispatch_priority(inner: &Arc<Inner>, item: (i64, QueueItemId)) {
193 let dispatch = claim_priority(&mut inner.queue.lock(), item);
194 match dispatch {
195 Dispatch::AlreadyRunning => {}
196 Dispatch::Requeued => {
197 inner.has_work.notify_one();
198 }
199 Dispatch::Spawn => {
200 let spawn_inner = inner.clone();
201 let spawned = std::thread::Builder::new()
202 .name("koan-dl-prio".into())
203 .spawn(move || {
204 let _claim = Claim {
205 inner: spawn_inner.clone(),
206 db_id: item.0,
207 priority: true,
208 };
209 run_download(&spawn_inner, item);
210 });
211 if let Err(e) = spawned {
212 log::error!("failed to spawn priority download: {}", e);
213 let mut q = inner.queue.lock();
214 release_priority(&mut q, item.0);
215 q.pending.push_front(item);
216 drop(q);
217 inner.has_work.notify_one();
218 }
219 }
220 }
221}
222
223fn settle_waiters(inner: &Arc<Inner>, db_id: i64, downloaded: QueueItemId) {
229 let waiting: Vec<QueueItemId> = {
230 let q = inner.queue.lock();
231 q.in_flight
232 .get(&db_id)
233 .map(|ids| ids.iter().copied().filter(|id| *id != downloaded).collect())
234 .unwrap_or_default()
235 };
236 if waiting.is_empty() {
237 return;
238 }
239 let Some(item) = inner.state.get_item(downloaded) else {
240 return;
241 };
242 for id in waiting {
243 inner.state.update_paths(&[(id, item.path.clone())]);
244 inner.state.update_item_state(id, item.state.clone());
245 if inner.state.is_cursor(id) {
248 inner.cmd_tx.send(PlayerCommand::TrackReady(id)).ok();
249 }
250 }
251}
252
253fn run_download(inner: &Arc<Inner>, (db_id, queue_id): (i64, QueueItemId)) {
255 let Some(client) = inner.client.as_ref() else {
256 crate::helpers::fail_track(
259 &inner.state,
260 &inner.cmd_tx,
261 queue_id,
262 crate::helpers::remote_unavailable(&inner.cfg),
263 );
264 return;
265 };
266
267 let outcome = std::panic::catch_unwind(AssertUnwindSafe(|| {
268 download_track(
269 db_id,
270 queue_id,
271 &inner.cmd_tx,
272 &inner.log_buf,
273 &inner.state,
274 &inner.cfg,
275 client,
276 );
277 }));
278
279 if outcome.is_err() {
280 log::error!("download panicked for {:?}", queue_id);
281 crate::helpers::fail_track(
282 &inner.state,
283 &inner.cmd_tx,
284 queue_id,
285 "download panicked".into(),
286 );
287 }
288
289 settle_waiters(inner, db_id, queue_id);
292}
293
294fn worker_loop(inner: Arc<Inner>) {
296 loop {
297 let item = {
298 let mut q = inner.queue.lock();
299 loop {
300 match q.pending.pop_front() {
301 Some(item) => {
302 match q.in_flight.get_mut(&item.0) {
305 Some(waiting) => {
306 waiting.insert(item.1);
307 }
308 None => {
309 q.in_flight.insert(item.0, HashSet::from([item.1]));
310 break item;
311 }
312 }
313 }
314 None => inner.has_work.wait(&mut q),
315 }
316 }
317 };
318 let _claim = Claim {
319 inner: inner.clone(),
320 db_id: item.0,
321 priority: false,
322 };
323 run_download(&inner, item);
324 }
325}
326
327fn cursor_watcher(inner: Arc<Inner>) {
330 let changed = crate::signal::engine_changed();
331 let mut seen = changed.generation();
332 let mut last_cursor: Option<QueueItemId> = None;
333 loop {
334 seen = changed.wait(seen);
335
336 let current = inner.state.cursor();
337 if current == last_cursor {
338 continue;
339 }
340 last_cursor = current;
341
342 let Some(cursor_id) = current else {
343 continue;
344 };
345
346 let is_pending = inner
347 .state
348 .item_load_state(cursor_id)
349 .is_some_and(|s| matches!(s, LoadState::Pending));
350 if !is_pending {
351 continue;
352 }
353
354 let album_mate_ids: HashSet<QueueItemId> = inner
355 .state
356 .same_album_item_ids(cursor_id)
357 .into_iter()
358 .collect();
359
360 let mut priority_items = Vec::new();
361 {
362 let mut q = inner.queue.lock();
363 if let Some(pos) = q.pending.iter().position(|(_, qid)| *qid == cursor_id) {
364 priority_items.push(q.pending.remove(pos).expect("position just found"));
365
366 if !album_mate_ids.is_empty() {
367 bump_to_front(&mut q.pending, &album_mate_ids);
368 }
369
370 if let Some(next) = q.pending.pop_front() {
372 priority_items.push(next);
373 }
374 }
375 }
376
377 for item in priority_items {
378 dispatch_priority(&inner, item);
379 }
380 }
381}
382
383pub fn shared(
393 cmd_tx: &crossbeam_channel::Sender<PlayerCommand>,
394 state: &Arc<SharedPlayerState>,
395 log_buf: Option<Arc<StdMutex<Vec<String>>>>,
396) -> &'static DownloadQueue {
397 static QUEUE: std::sync::OnceLock<DownloadQueue> = std::sync::OnceLock::new();
398 QUEUE.get_or_init(|| {
399 DownloadQueue::spawn(
400 cmd_tx.clone(),
401 state.clone(),
402 log_buf.unwrap_or_else(|| Arc::new(StdMutex::new(Vec::new()))),
403 )
404 })
405}
406
407#[cfg(test)]
408mod tests {
409 use super::*;
410
411 fn qid() -> QueueItemId {
412 QueueItemId::new()
413 }
414
415 #[test]
416 fn priority_lane_never_exceeds_its_permits() {
417 let mut q = Queue::default();
418
419 let mut spawned = 0;
421 for i in 0..500 {
422 if claim_priority(&mut q, (i, qid())) == Dispatch::Spawn {
423 spawned += 1;
424 }
425 assert!(
426 q.priority_active <= PRIORITY_PERMITS,
427 "priority lane over its permit count at iteration {}",
428 i
429 );
430 }
431
432 assert_eq!(spawned, PRIORITY_PERMITS, "only permitted claims may spawn");
433 assert_eq!(
434 q.pending.len(),
435 500 - PRIORITY_PERMITS,
436 "everything else must be queued, not dropped"
437 );
438 }
439
440 #[test]
441 fn released_permits_are_reusable() {
442 let mut q = Queue::default();
443 assert_eq!(claim_priority(&mut q, (1, qid())), Dispatch::Spawn);
444 assert_eq!(claim_priority(&mut q, (2, qid())), Dispatch::Spawn);
445 assert_eq!(claim_priority(&mut q, (3, qid())), Dispatch::Requeued);
446
447 release_priority(&mut q, 1);
448 assert_eq!(claim_priority(&mut q, (4, qid())), Dispatch::Spawn);
449 assert!(q.priority_active <= PRIORITY_PERMITS);
450 }
451
452 #[test]
453 fn an_in_flight_track_is_never_claimed_twice() {
454 let mut q = Queue::default();
455 let id = qid();
456 assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::Spawn);
457 assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::AlreadyRunning);
458 assert_eq!(q.priority_active, 1);
459 assert!(
460 q.pending.is_empty(),
461 "a duplicate request must not re-queue the track"
462 );
463 }
464
465 #[test]
466 fn playing_a_track_again_joins_the_transfer_already_running() {
467 let mut q = Queue::default();
473 let (first, again) = (qid(), qid());
474 assert_eq!(claim_priority(&mut q, (7, first)), Dispatch::Spawn);
475 assert_eq!(claim_priority(&mut q, (7, again)), Dispatch::AlreadyRunning);
476
477 assert_eq!(q.priority_active, 1, "one transfer, not two");
478 assert!(q.pending.is_empty());
479 assert_eq!(
480 q.in_flight.get(&7),
481 Some(&HashSet::from([first, again])),
482 "both entries wait on the one transfer"
483 );
484 }
485
486 #[test]
487 fn a_worker_picking_up_a_duplicate_waits_on_the_running_one() {
488 let mut q = Queue::default();
490 let (running, queued) = (qid(), qid());
491 assert_eq!(claim_priority(&mut q, (7, running)), Dispatch::Spawn);
492
493 match q.in_flight.get_mut(&7) {
495 Some(waiting) => {
496 waiting.insert(queued);
497 }
498 None => panic!("the track should already be claimed"),
499 }
500
501 assert_eq!(
502 q.in_flight.get(&7),
503 Some(&HashSet::from([running, queued])),
504 "the queued entry waits rather than starting a second transfer"
505 );
506 }
507
508 #[test]
509 fn different_tracks_still_run_side_by_side() {
510 let mut q = Queue::default();
512 assert_eq!(claim_priority(&mut q, (1, qid())), Dispatch::Spawn);
513 assert_eq!(claim_priority(&mut q, (2, qid())), Dispatch::Spawn);
514 assert_eq!(q.priority_active, 2);
515 }
516
517 #[test]
518 fn requeued_priority_item_goes_to_the_head_of_the_queue() {
519 let mut q = Queue::default();
520 q.pending.push_back((9, qid()));
521 for i in 0..PRIORITY_PERMITS {
522 claim_priority(&mut q, (i as i64, qid()));
523 }
524
525 let wanted = qid();
526 assert_eq!(claim_priority(&mut q, (7, wanted)), Dispatch::Requeued);
527 assert_eq!(q.pending.front().map(|(_, id)| *id), Some(wanted));
528 }
529
530 #[test]
531 fn claiming_removes_a_duplicate_queue_entry() {
532 let mut q = Queue::default();
533 let id = qid();
534 q.pending.push_back((1, id));
535 q.pending.push_back((2, qid()));
536
537 assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::Spawn);
538 assert_eq!(
539 q.pending.len(),
540 1,
541 "the pool must not also pick up the claimed track"
542 );
543 }
544
545 #[test]
546 fn bump_to_front_preserves_relative_order() {
547 let (a, b, c, d) = (qid(), qid(), qid(), qid());
548 let mut pending: VecDeque<(i64, QueueItemId)> =
549 [(1, a), (2, b), (3, c), (4, d)].into_iter().collect();
550 let mates: HashSet<QueueItemId> = [b, d].into_iter().collect();
551
552 bump_to_front(&mut pending, &mates);
553
554 let order: Vec<QueueItemId> = pending.iter().map(|(_, id)| *id).collect();
555 assert_eq!(order, vec![b, d, a, c]);
556 }
557}