Skip to main content

koan_core/remote/
queue.rs

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
14/// Concurrent downloads the priority lane may run outside the worker pool.
15/// Small on purpose: its job is to get the track under the cursor playing, and
16/// every extra request competes with it for the same link.
17const PRIORITY_PERMITS: usize = 2;
18
19/// Persistent download queue — lives for the app's lifetime.
20///
21/// Items are submitted via `enqueue()` and downloaded by a fixed pool of worker
22/// threads. Cursor changes reorder the queue so the current track downloads
23/// first, followed by same-album tracks for gapless playback; those jump the
24/// queue through a permit-limited priority lane rather than by spawning
25/// unbounded threads.
26#[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    /// `None` when remote is not configured — nothing is downloadable.
39    client: Option<Arc<SubsonicClient>>,
40}
41
42/// Queue state and the in-flight bookkeeping that keeps a track from being
43/// downloaded by two threads at once.
44#[derive(Default)]
45struct Queue {
46    pending: VecDeque<(i64, QueueItemId)>,
47    /// Tracks being fetched, and every queue entry waiting on each.
48    ///
49    /// Keyed by track, because the track decides which file the download
50    /// writes. Keyed by queue entry it did not dedupe anything that mattered:
51    /// playing something a second time before it had arrived made a new entry
52    /// with a new id, so nothing matched and a second transfer started over
53    /// the first — two threads truncating and writing one `.part`, and
54    /// whichever finished first renaming it out from under the other.
55    in_flight: HashMap<i64, HashSet<QueueItemId>>,
56    priority_active: usize,
57}
58
59/// What a priority request should do, given the state of the lane.
60#[derive(Debug, PartialEq, Eq)]
61enum Dispatch {
62    /// A permit was taken and the item claimed — spawn a thread for it.
63    Spawn,
64    /// No permit free; the item now sits at the head of the work queue.
65    Requeued,
66    /// Some thread is already downloading it.
67    AlreadyRunning,
68}
69
70/// Claim `item` for the priority lane, or push it to the front of the queue if
71/// every permit is taken. On `Spawn` the caller owns the claim and must release
72/// it via `release_priority` when the download ends.
73fn 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    // Already being fetched: wait on it rather than fetch it again. The entry
78    // is remembered so it gets the answer when the one transfer lands.
79    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
97/// Releases an in-flight claim however the download ends — including a panic.
98struct 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    /// Spawn the download queue with persistent worker threads.
117    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    /// Add items to the download queue.
161    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    /// Submit a single item for priority download (e.g. user clicked a Pending
170    /// track). Also bumps same-album pending tracks for gapless playback.
171    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
183/// Move every item whose id is in `ids` ahead of the rest, preserving order.
184fn 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
191/// Start a priority download, or queue it at the front when the lane is full.
192fn 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
223/// Hand every other entry waiting on this track the answer the transfer got.
224///
225/// Two entries for one track are one download and two things to tell. Without
226/// this the second sits `Pending` forever, waiting for a transfer that already
227/// finished and will not run again.
228fn 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        // The player only wakes for this, so an entry the cursor is sitting on
246        // would otherwise wait on a download that has already happened.
247        if inner.state.is_cursor(id) {
248            inner.cmd_tx.send(PlayerCommand::TrackReady(id)).ok();
249        }
250    }
251}
252
253/// Run one download, containing any panic so the worker pool never shrinks.
254fn run_download(inner: &Arc<Inner>, (db_id, queue_id): (i64, QueueItemId)) {
255    let Some(client) = inner.client.as_ref() else {
256        // Failed, not left Pending: the player waits for Ready, so a queue of
257        // tracks that can never arrive would otherwise sit saying nothing.
258        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    // Before the claim is released, while the waiting entries are still
290    // recorded against this track.
291    settle_waiters(inner, db_id, queue_id);
292}
293
294/// Worker loop: wait for work, download, repeat.
295fn 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                        // Already being fetched: this entry waits on the one
303                        // transfer rather than starting a second over it.
304                        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
327/// Cursor watcher: when the cursor moves to a pending track, hand it and the
328/// next track to the priority lane and bump same-album tracks to the front.
329fn 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                // Grab the next track too, for gapless lookahead.
371                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
383/// The process's download queue.
384///
385/// One player means one pool, one priority lane and one cursor watcher; a
386/// second set would compete with the first for the same link and the same
387/// cursor. Every front end reaches downloads through here — the TUI directly,
388/// the FFI and the GraphQL server through `helpers::spawn_downloads`.
389///
390/// `log_buf` is only honoured by whoever initialises it, which is the TUI when
391/// it is running, since it is the only front end that shows the buffer.
392pub 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        // Rapid cursor movement: a fresh track lands on the lane every poll.
420        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        // Playing something twice before it has arrived makes a second queue
468        // entry with an id of its own. The track is the same, and so is the
469        // file a download would write — two of them would truncate and write
470        // over one another, and whichever finished first would rename it away
471        // from the other.
472        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        // The same, arriving through the queue rather than the priority lane.
489        let mut q = Queue::default();
490        let (running, queued) = (qid(), qid());
491        assert_eq!(claim_priority(&mut q, (7, running)), Dispatch::Spawn);
492
493        // What `worker_loop` does with the next pending item.
494        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        // Keying on the track must not serialise unrelated downloads.
511        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}