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    /// When the cache was last trimmed to its limit.
41    last_evicted: Mutex<Option<std::time::Instant>>,
42}
43
44/// Queue state and the in-flight bookkeeping that keeps a track from being
45/// downloaded by two threads at once.
46#[derive(Default)]
47struct Queue {
48    pending: VecDeque<(i64, QueueItemId)>,
49    /// Tracks being fetched, and every queue entry waiting on each.
50    ///
51    /// Keyed by track, because the track decides which file the download
52    /// writes. Keyed by queue entry it did not dedupe anything that mattered:
53    /// playing something a second time before it had arrived made a new entry
54    /// with a new id, so nothing matched and a second transfer started over
55    /// the first — two threads truncating and writing one `.part`, and
56    /// whichever finished first renaming it out from under the other.
57    in_flight: HashMap<i64, HashSet<QueueItemId>>,
58    priority_active: usize,
59}
60
61/// What a priority request should do, given the state of the lane.
62#[derive(Debug, PartialEq, Eq)]
63enum Dispatch {
64    /// A permit was taken and the item claimed — spawn a thread for it.
65    Spawn,
66    /// No permit free; the item now sits at the head of the work queue.
67    Requeued,
68    /// Some thread is already downloading it.
69    AlreadyRunning,
70}
71
72/// Claim `item` for the priority lane, or push it to the front of the queue if
73/// every permit is taken. On `Spawn` the caller owns the claim and must release
74/// it via `release_priority` when the download ends.
75fn 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    // Already being fetched: wait on it rather than fetch it again. The entry
80    // is remembered so it gets the answer when the one transfer lands.
81    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
99/// Releases an in-flight claim however the download ends — including a panic.
100struct 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    /// Spawn the download queue with persistent worker threads.
119    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    /// Add items to the download queue.
169    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    /// Submit a single item for priority download (e.g. user clicked a Pending
178    /// track). Also bumps same-album pending tracks for gapless playback.
179    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
191/// Move every item whose id is in `ids` ahead of the rest, preserving order.
192fn 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
199/// Start a priority download, or queue it at the front when the lane is full.
200fn 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
231/// Hand every other entry waiting on this track the answer the transfer got.
232///
233/// Two entries for one track are one download and two things to tell. Without
234/// this the second sits `Pending` forever, waiting for a transfer that already
235/// finished and will not run again.
236fn 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        // The player only wakes for this, so an entry the cursor is sitting on
254        // would otherwise wait on a download that has already happened.
255        if inner.state.is_cursor(id) {
256            inner.cmd_tx.send(PlayerCommand::TrackReady(id)).ok();
257        }
258    }
259}
260
261/// Run one download, containing any panic so the worker pool never shrinks.
262fn run_download(inner: &Arc<Inner>, (db_id, queue_id): (i64, QueueItemId)) {
263    let Some(client) = inner.client.as_ref() else {
264        // Failed, not left Pending: the player waits for Ready, so a queue of
265        // tracks that can never arrive would otherwise sit saying nothing.
266        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    // Before the claim is released, while the waiting entries are still
298    // recorded against this track.
299    settle_waiters(inner, db_id, queue_id);
300    trim_cache(inner);
301}
302
303/// How often the cache is checked against its limit, at most: each download
304/// adds to it, and a check reads the whole cache's size from the database.
305const EVICT_EVERY: std::time::Duration = std::time::Duration::from_secs(60);
306
307/// Trim the cache to its configured limit, keeping everything in the queue:
308/// the player may be reading those files, and they are what is wanted next.
309/// The limit is read afresh, so one set in Settings applies without a
310/// restart.
311fn 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
338/// Worker loop: wait for work, download, repeat.
339///
340/// While the server is down, workers take nothing new: the download that found
341/// it down waits it out, and everything behind it stays `Pending` rather than
342/// each piling onto a server that is not answering.
343fn 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                    // Already being fetched: this entry waits on the one
353                    // transfer rather than starting a second over it.
354                    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
376/// Cursor watcher: when the cursor moves to a pending track, hand it and the
377/// next track to the priority lane and bump same-album tracks to the front.
378fn 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                // Grab the next track too, for gapless lookahead.
420                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
432/// The process's download queue.
433///
434/// One player means one pool, one priority lane and one cursor watcher; a
435/// second set would compete with the first for the same link and the same
436/// cursor. Every front end reaches downloads through here — the TUI directly,
437/// the FFI and the GraphQL server through `helpers::spawn_downloads`.
438///
439/// `log_buf` is only honoured by whoever initialises it, which is the TUI when
440/// it is running, since it is the only front end that shows the buffer.
441pub 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        // Rapid cursor movement: a fresh track lands on the lane every poll.
469        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        // Playing something twice before it has arrived makes a second queue
517        // entry with an id of its own. The track is the same, and so is the
518        // file a download would write — two of them would truncate and write
519        // over one another, and whichever finished first would rename it away
520        // from the other.
521        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        // The same, arriving through the queue rather than the priority lane.
538        let mut q = Queue::default();
539        let (running, queued) = (qid(), qid());
540        assert_eq!(claim_priority(&mut q, (7, running)), Dispatch::Spawn);
541
542        // What `worker_loop` does with the next pending item.
543        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        // Keying on the track must not serialise unrelated downloads.
560        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}