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};
10
11use crate::helpers::download_track;
12
13/// Concurrent downloads the priority lane may run outside the worker pool.
14/// Small on purpose: its job is to get the track under the cursor playing, and
15/// every extra request competes with it for the same link.
16const PRIORITY_PERMITS: usize = 2;
17
18/// Persistent download queue — lives for the app's lifetime.
19///
20/// Items are submitted via `enqueue()` and downloaded by a fixed pool of worker
21/// threads. Cursor changes reorder the queue so the current track downloads
22/// first, followed by same-album tracks for gapless playback; those jump the
23/// queue through a permit-limited priority lane rather than by spawning
24/// unbounded threads.
25#[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    /// When the cache was last trimmed to its limit.
37    last_evicted: Mutex<Option<std::time::Instant>>,
38    /// Worker threads started so far. Grows to the configured count as work
39    /// arrives; a worker past the current count parks rather than exits.
40    spawned: std::sync::atomic::AtomicUsize,
41}
42
43/// The parallel-downloads setting as it stands now, not as it stood when the
44/// queue was made: the queue lives as long as the process.
45fn workers_allowed() -> usize {
46    config::Config::cached()
47        .remote
48        .download_workers
49        .clamp(1, 16)
50}
51
52/// Start workers until there are as many as the setting allows.
53fn 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
80/// What a worker should fetch next: the track under the cursor ahead of
81/// everything, nothing at all while that track is already being fetched, and
82/// otherwise the front of the queue.
83fn 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
96/// Someone asked for music: try the server now rather than when the outage
97/// backoff next says to. The backoff is for retries nobody is waiting on; a
98/// minute of it, earned while the phone was in a pocket with no signal, is not
99/// how long a tap on play should take.
100fn retry_server_now() {
101    if let Some(client) = crate::helpers::subsonic_client(&config::Config::cached()) {
102        client.outage().retry_now();
103    }
104}
105
106/// The track under the cursor, when it still has to be fetched.
107fn cursor_download(state: &SharedPlayerState) -> Option<(i64, QueueItemId)> {
108    let id = state.cursor()?;
109    let item = state.get_item(id)?;
110    if !matches!(item.state, crate::player::state::ItemState::Pending) {
111        return None;
112    }
113    Some((item.db_id?, id))
114}
115
116/// Queue state and the in-flight bookkeeping that keeps a track from being
117/// downloaded by two threads at once.
118#[derive(Default)]
119struct Queue {
120    pending: VecDeque<(i64, QueueItemId)>,
121    /// Tracks being fetched, and every queue entry waiting on each.
122    ///
123    /// Keyed by track, because the track decides which file the download
124    /// writes. Keyed by queue entry it would dedupe nothing that matters:
125    /// playing something a second time before it has arrived makes a new entry
126    /// with a new id, so nothing would match and a second transfer would start
127    /// over the first — two threads truncating and writing one `.part`, and
128    /// whichever finishes first renaming it out from under the other.
129    in_flight: HashMap<i64, HashSet<QueueItemId>>,
130    priority_active: usize,
131}
132
133/// What a priority request should do, given the state of the lane.
134#[derive(Debug, PartialEq, Eq)]
135enum Dispatch {
136    /// A permit was taken and the item claimed — spawn a thread for it.
137    Spawn,
138    /// No permit free; the item now sits at the head of the work queue.
139    Requeued,
140    /// Some thread is already downloading it.
141    AlreadyRunning,
142}
143
144/// Claim `item` for the priority lane, or push it to the front of the queue if
145/// every permit is taken. On `Spawn` the caller owns the claim and must release
146/// it via `release_priority` when the download ends.
147fn claim_priority(q: &mut Queue, item: (i64, QueueItemId)) -> Dispatch {
148    let (db_id, queue_id) = item;
149    q.pending.retain(|(_, qid)| *qid != queue_id);
150
151    // Already being fetched: wait on it rather than fetch it again. The entry
152    // is remembered so it gets the answer when the one transfer lands.
153    if let Some(waiting) = q.in_flight.get_mut(&db_id) {
154        waiting.insert(queue_id);
155        return Dispatch::AlreadyRunning;
156    }
157    if q.priority_active >= PRIORITY_PERMITS {
158        q.pending.push_front(item);
159        return Dispatch::Requeued;
160    }
161    q.priority_active += 1;
162    q.in_flight.insert(db_id, HashSet::from([queue_id]));
163    Dispatch::Spawn
164}
165
166fn release_priority(q: &mut Queue, db_id: i64) {
167    q.in_flight.remove(&db_id);
168    q.priority_active = q.priority_active.saturating_sub(1);
169}
170
171/// Releases an in-flight claim however the download ends — including a panic.
172struct Claim {
173    inner: Arc<Inner>,
174    db_id: i64,
175    priority: bool,
176}
177
178impl Drop for Claim {
179    fn drop(&mut self) {
180        let mut q = self.inner.queue.lock();
181        if self.priority {
182            release_priority(&mut q, self.db_id);
183        } else {
184            q.in_flight.remove(&self.db_id);
185        }
186        drop(q);
187        // Workers held back for the track under the cursor go on from here.
188        self.inner.has_work.notify_all();
189    }
190}
191
192impl DownloadQueue {
193    /// Spawn the download queue with persistent worker threads.
194    pub fn spawn(
195        cmd_tx: crossbeam_channel::Sender<PlayerCommand>,
196        state: Arc<SharedPlayerState>,
197        log_buf: Arc<StdMutex<Vec<String>>>,
198    ) -> Self {
199        let inner = Arc::new(Inner {
200            queue: Mutex::new(Queue::default()),
201            has_work: Condvar::new(),
202            state,
203            cmd_tx,
204            log_buf,
205            last_evicted: Mutex::new(None),
206            spawned: std::sync::atomic::AtomicUsize::new(0),
207        });
208        ensure_workers(&inner);
209
210        let trimmer = inner.clone();
211        let _ = std::thread::Builder::new()
212            .name("koan-dl-trim".into())
213            .spawn(move || trim_cache(&trimmer));
214
215        let watcher_inner = inner.clone();
216        if let Err(e) = std::thread::Builder::new()
217            .name("koan-dl-watch".into())
218            .spawn(move || cursor_watcher(watcher_inner))
219        {
220            log::error!("failed to spawn download cursor watcher: {}", e);
221        }
222
223        Self { inner }
224    }
225
226    /// Add items to the download queue.
227    pub fn enqueue(&self, items: Vec<(i64, QueueItemId)>) {
228        if items.is_empty() {
229            return;
230        }
231        ensure_workers(&self.inner);
232        retry_server_now();
233        self.inner.queue.lock().pending.extend(items);
234        self.inner.has_work.notify_all();
235        if let Some(cursor) = self.inner.state.cursor() {
236            promote_cursor(&self.inner, cursor);
237        }
238    }
239
240    /// Submit a single item for priority download (e.g. user clicked a Pending
241    /// track). Also bumps same-album pending tracks for gapless playback.
242    pub fn prioritize(&self, db_id: i64, queue_id: QueueItemId) {
243        retry_server_now();
244        dispatch_priority(&self.inner, (db_id, queue_id));
245
246        let album_mates = self.inner.state.same_album_item_ids(queue_id);
247        if !album_mates.is_empty() {
248            let mate_set: HashSet<QueueItemId> = album_mates.into_iter().collect();
249            bump_to_front(&mut self.inner.queue.lock().pending, &mate_set);
250            self.inner.has_work.notify_all();
251        }
252    }
253}
254
255/// Move every item whose id is in `ids` ahead of the rest, preserving order.
256fn bump_to_front(pending: &mut VecDeque<(i64, QueueItemId)>, ids: &HashSet<QueueItemId>) {
257    let (front, rest): (VecDeque<_>, VecDeque<_>) =
258        pending.drain(..).partition(|(_, qid)| ids.contains(qid));
259    *pending = front;
260    pending.extend(rest);
261}
262
263/// Start a priority download, or queue it at the front when the lane is full.
264fn dispatch_priority(inner: &Arc<Inner>, item: (i64, QueueItemId)) {
265    let dispatch = claim_priority(&mut inner.queue.lock(), item);
266    match dispatch {
267        Dispatch::AlreadyRunning => {}
268        Dispatch::Requeued => {
269            inner.has_work.notify_one();
270        }
271        Dispatch::Spawn => {
272            let spawn_inner = inner.clone();
273            let spawned = std::thread::Builder::new()
274                .name("koan-dl-prio".into())
275                .spawn(move || {
276                    let _claim = Claim {
277                        inner: spawn_inner.clone(),
278                        db_id: item.0,
279                        priority: true,
280                    };
281                    run_download(&spawn_inner, item);
282                });
283            if let Err(e) = spawned {
284                log::error!("failed to spawn priority download: {}", e);
285                let mut q = inner.queue.lock();
286                release_priority(&mut q, item.0);
287                q.pending.push_front(item);
288                drop(q);
289                inner.has_work.notify_one();
290            }
291        }
292    }
293}
294
295/// Hand every other entry waiting on this track the answer the transfer got.
296///
297/// Two entries for one track are one download and two things to tell. Without
298/// this the second sits `Pending` forever, waiting for a transfer that already
299/// finished and will not run again.
300fn settle_waiters(inner: &Arc<Inner>, db_id: i64, downloaded: QueueItemId) {
301    let waiting: Vec<QueueItemId> = {
302        let q = inner.queue.lock();
303        q.in_flight
304            .get(&db_id)
305            .map(|ids| ids.iter().copied().filter(|id| *id != downloaded).collect())
306            .unwrap_or_default()
307    };
308    if waiting.is_empty() {
309        return;
310    }
311    let Some(item) = inner.state.get_item(downloaded) else {
312        return;
313    };
314    for id in waiting {
315        inner.state.update_paths(&[(id, item.path.clone())]);
316        inner.state.update_item_state(id, item.state.clone());
317        // The player only wakes for this, so an entry the cursor is sitting on
318        // would otherwise wait on a download that has already happened.
319        if inner.state.is_cursor(id) {
320            inner.cmd_tx.send(PlayerCommand::TrackReady(id)).ok();
321        }
322    }
323}
324
325/// Run one download, containing any panic so the worker pool never shrinks.
326///
327/// The client is looked up per download, not once for the queue's lifetime:
328/// the queue lives as long as the process, and signing in, out or elsewhere
329/// has to reach it.
330fn run_download(inner: &Arc<Inner>, (db_id, queue_id): (i64, QueueItemId)) {
331    let cfg = config::Config::cached();
332    let Some(client) = crate::helpers::subsonic_client(&cfg) else {
333        // Failed, not left Pending: the player waits for Ready, so a queue of
334        // tracks that can never arrive would otherwise sit saying nothing.
335        crate::helpers::fail_track(
336            &inner.state,
337            &inner.cmd_tx,
338            queue_id,
339            crate::helpers::remote_unavailable(&cfg),
340        );
341        return;
342    };
343
344    let outcome = std::panic::catch_unwind(AssertUnwindSafe(|| {
345        download_track(
346            db_id,
347            queue_id,
348            &inner.cmd_tx,
349            &inner.log_buf,
350            &inner.state,
351            &cfg,
352            &client,
353        );
354    }));
355
356    if outcome.is_err() {
357        log::error!("download panicked for {:?}", queue_id);
358        crate::helpers::fail_track(
359            &inner.state,
360            &inner.cmd_tx,
361            queue_id,
362            "download panicked".into(),
363        );
364    }
365
366    // Before the claim is released, while the waiting entries are still
367    // recorded against this track.
368    settle_waiters(inner, db_id, queue_id);
369    trim_cache(inner);
370}
371
372/// How often the cache is checked against its limit, at most: each download
373/// adds to it, and a check reads the whole cache's size from the database.
374const EVICT_EVERY: std::time::Duration = std::time::Duration::from_secs(60);
375
376/// Trim the cache to its configured limit, keeping everything in the queue:
377/// the player may be reading those files, and they are what is wanted next.
378/// The limit is read afresh, so one set in Settings applies without a
379/// restart.
380fn trim_cache(inner: &Inner) {
381    {
382        let mut last = inner.last_evicted.lock();
383        if last.is_some_and(|t| t.elapsed() < EVICT_EVERY) {
384            return;
385        }
386        *last = Some(std::time::Instant::now());
387    }
388    let cfg = config::Config::cached();
389    if cfg.cache_limit_bytes().is_none() {
390        return;
391    }
392    let keep = inner
393        .state
394        .snapshot_playlist()
395        .0
396        .iter()
397        .filter_map(|i| i.db_id)
398        .collect();
399    match crate::db::pool::shared().get() {
400        Ok(db) => {
401            crate::helpers::evict_cache(&db, &cfg, &keep, false);
402        }
403        Err(e) => log::warn!("cache eviction: could not open the database: {e}"),
404    }
405}
406
407/// Worker loop: wait for work, download, repeat.
408///
409/// While the server is down, workers take nothing new: the download that found
410/// it down waits it out, and everything behind it stays `Pending` rather than
411/// each piling onto a server that is not answering.
412/// One download worker. `index` is its place in the pool: a worker past the
413/// current parallel-downloads setting parks until the setting lets it work.
414///
415/// While the track under the cursor is still being fetched, no worker starts
416/// another transfer. On a slow link every parallel download takes a share of
417/// the bandwidth, and the one track someone is waiting to hear would arrive
418/// last among equals; the rest of the album can follow it.
419fn worker_loop(inner: Arc<Inner>, index: usize) {
420    loop {
421        let item = loop {
422            if let Some(client) = crate::helpers::subsonic_client(&config::Config::cached()) {
423                client.outage().hold();
424            }
425            // Read before the queue lock, which is never held across the
426            // player state's.
427            let cursor = cursor_download(&inner.state);
428            let mut q = inner.queue.lock();
429            if index >= workers_allowed() {
430                inner.has_work.wait(&mut q);
431                continue;
432            }
433            match next_item(&mut q, cursor) {
434                Some(item) => {
435                    // Already being fetched: this entry waits on the one
436                    // transfer rather than starting a second over it.
437                    match q.in_flight.get_mut(&item.0) {
438                        Some(waiting) => {
439                            waiting.insert(item.1);
440                        }
441                        None => {
442                            q.in_flight.insert(item.0, HashSet::from([item.1]));
443                            break item;
444                        }
445                    }
446                }
447                None => inner.has_work.wait(&mut q),
448            }
449        };
450        let _claim = Claim {
451            inner: inner.clone(),
452            db_id: item.0,
453            priority: false,
454        };
455        run_download(&inner, item);
456    }
457}
458
459/// Cursor watcher: when the cursor moves to a pending track, hand it and the
460/// next track to the priority lane and bump same-album tracks to the front.
461///
462/// Also wakes the workers, which hold back while the cursor's track is being
463/// fetched: one that looked before the cursor moved is waiting on a track
464/// nobody wants any more.
465fn cursor_watcher(inner: Arc<Inner>) {
466    let changed = crate::signal::engine_changed();
467    let mut seen = changed.generation();
468    let mut last_cursor: Option<QueueItemId> = None;
469    loop {
470        seen = changed.wait(seen);
471
472        let current = inner.state.cursor();
473        if current == last_cursor {
474            continue;
475        }
476        last_cursor = current;
477        inner.has_work.notify_all();
478
479        if let Some(cursor_id) = current {
480            promote_cursor(&inner, cursor_id);
481        }
482    }
483}
484
485/// Send the cursor's track, and the one after it, down the priority lane, if
486/// the cursor's track is waiting in the queue.
487///
488/// Called when the cursor moves and when tracks are queued, because either can
489/// happen first: playing an album sends the new queue to the player and queues
490/// its downloads at once, and the player may set the cursor before or after.
491/// Waiting for the cursor alone missed the first play after launch, when the
492/// queue (and this watcher) did not exist until those downloads made it.
493fn promote_cursor(inner: &Arc<Inner>, cursor_id: QueueItemId) {
494    let is_pending = inner
495        .state
496        .item_load_state(cursor_id)
497        .is_some_and(|s| matches!(s, LoadState::Pending));
498    if !is_pending {
499        return;
500    }
501
502    let album_mate_ids: HashSet<QueueItemId> = inner
503        .state
504        .same_album_item_ids(cursor_id)
505        .into_iter()
506        .collect();
507
508    let mut priority_items = Vec::new();
509    {
510        let mut q = inner.queue.lock();
511        if let Some(pos) = q.pending.iter().position(|(_, qid)| *qid == cursor_id) {
512            priority_items.push(q.pending.remove(pos).expect("position just found"));
513
514            if !album_mate_ids.is_empty() {
515                bump_to_front(&mut q.pending, &album_mate_ids);
516            }
517
518            // Grab the next track too, for gapless lookahead.
519            if let Some(next) = q.pending.pop_front() {
520                priority_items.push(next);
521            }
522        }
523    }
524
525    for item in priority_items {
526        dispatch_priority(inner, item);
527    }
528}
529
530/// The process's download queue.
531///
532/// One player means one pool, one priority lane and one cursor watcher; a
533/// second set would compete with the first for the same link and the same
534/// cursor. Every front end reaches downloads through here — the TUI directly,
535/// the FFI and the GraphQL server through `helpers::spawn_downloads`.
536///
537/// `log_buf` is only honoured by whoever initialises it, which is the TUI when
538/// it is running, since it is the only front end that shows the buffer.
539pub fn shared(
540    cmd_tx: &crossbeam_channel::Sender<PlayerCommand>,
541    state: &Arc<SharedPlayerState>,
542    log_buf: Option<Arc<StdMutex<Vec<String>>>>,
543) -> &'static DownloadQueue {
544    static QUEUE: std::sync::OnceLock<DownloadQueue> = std::sync::OnceLock::new();
545    QUEUE.get_or_init(|| {
546        DownloadQueue::spawn(
547            cmd_tx.clone(),
548            state.clone(),
549            log_buf.unwrap_or_else(|| Arc::new(StdMutex::new(Vec::new()))),
550        )
551    })
552}
553
554#[cfg(test)]
555mod tests {
556    use super::*;
557
558    fn qid() -> QueueItemId {
559        QueueItemId::new()
560    }
561
562    #[test]
563    fn priority_lane_never_exceeds_its_permits() {
564        let mut q = Queue::default();
565
566        // Rapid cursor movement: a fresh track lands on the lane every poll.
567        let mut spawned = 0;
568        for i in 0..500 {
569            if claim_priority(&mut q, (i, qid())) == Dispatch::Spawn {
570                spawned += 1;
571            }
572            assert!(
573                q.priority_active <= PRIORITY_PERMITS,
574                "priority lane over its permit count at iteration {}",
575                i
576            );
577        }
578
579        assert_eq!(spawned, PRIORITY_PERMITS, "only permitted claims may spawn");
580        assert_eq!(
581            q.pending.len(),
582            500 - PRIORITY_PERMITS,
583            "everything else must be queued, not dropped"
584        );
585    }
586
587    #[test]
588    fn released_permits_are_reusable() {
589        let mut q = Queue::default();
590        assert_eq!(claim_priority(&mut q, (1, qid())), Dispatch::Spawn);
591        assert_eq!(claim_priority(&mut q, (2, qid())), Dispatch::Spawn);
592        assert_eq!(claim_priority(&mut q, (3, qid())), Dispatch::Requeued);
593
594        release_priority(&mut q, 1);
595        assert_eq!(claim_priority(&mut q, (4, qid())), Dispatch::Spawn);
596        assert!(q.priority_active <= PRIORITY_PERMITS);
597    }
598
599    #[test]
600    fn an_in_flight_track_is_never_claimed_twice() {
601        let mut q = Queue::default();
602        let id = qid();
603        assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::Spawn);
604        assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::AlreadyRunning);
605        assert_eq!(q.priority_active, 1);
606        assert!(
607            q.pending.is_empty(),
608            "a duplicate request must not re-queue the track"
609        );
610    }
611
612    #[test]
613    fn playing_a_track_again_joins_the_transfer_already_running() {
614        // Playing something twice before it has arrived makes a second queue
615        // entry with an id of its own. The track is the same, and so is the
616        // file a download would write — two of them would truncate and write
617        // over one another, and whichever finished first would rename it away
618        // from the other.
619        let mut q = Queue::default();
620        let (first, again) = (qid(), qid());
621        assert_eq!(claim_priority(&mut q, (7, first)), Dispatch::Spawn);
622        assert_eq!(claim_priority(&mut q, (7, again)), Dispatch::AlreadyRunning);
623
624        assert_eq!(q.priority_active, 1, "one transfer, not two");
625        assert!(q.pending.is_empty());
626        assert_eq!(
627            q.in_flight.get(&7),
628            Some(&HashSet::from([first, again])),
629            "both entries wait on the one transfer"
630        );
631    }
632
633    #[test]
634    fn a_worker_picking_up_a_duplicate_waits_on_the_running_one() {
635        // The same, arriving through the queue rather than the priority lane.
636        let mut q = Queue::default();
637        let (running, queued) = (qid(), qid());
638        assert_eq!(claim_priority(&mut q, (7, running)), Dispatch::Spawn);
639
640        // What `worker_loop` does with the next pending item.
641        match q.in_flight.get_mut(&7) {
642            Some(waiting) => {
643                waiting.insert(queued);
644            }
645            None => panic!("the track should already be claimed"),
646        }
647
648        assert_eq!(
649            q.in_flight.get(&7),
650            Some(&HashSet::from([running, queued])),
651            "the queued entry waits rather than starting a second transfer"
652        );
653    }
654
655    #[test]
656    fn different_tracks_still_run_side_by_side() {
657        // Keying on the track must not serialise unrelated downloads.
658        let mut q = Queue::default();
659        assert_eq!(claim_priority(&mut q, (1, qid())), Dispatch::Spawn);
660        assert_eq!(claim_priority(&mut q, (2, qid())), Dispatch::Spawn);
661        assert_eq!(q.priority_active, 2);
662    }
663
664    #[test]
665    fn requeued_priority_item_goes_to_the_head_of_the_queue() {
666        let mut q = Queue::default();
667        q.pending.push_back((9, qid()));
668        for i in 0..PRIORITY_PERMITS {
669            claim_priority(&mut q, (i as i64, qid()));
670        }
671
672        let wanted = qid();
673        assert_eq!(claim_priority(&mut q, (7, wanted)), Dispatch::Requeued);
674        assert_eq!(q.pending.front().map(|(_, id)| *id), Some(wanted));
675    }
676
677    #[test]
678    fn claiming_removes_a_duplicate_queue_entry() {
679        let mut q = Queue::default();
680        let id = qid();
681        q.pending.push_back((1, id));
682        q.pending.push_back((2, qid()));
683
684        assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::Spawn);
685        assert_eq!(
686            q.pending.len(),
687            1,
688            "the pool must not also pick up the claimed track"
689        );
690    }
691
692    #[test]
693    fn bump_to_front_preserves_relative_order() {
694        let (a, b, c, d) = (qid(), qid(), qid(), qid());
695        let mut pending: VecDeque<(i64, QueueItemId)> =
696            [(1, a), (2, b), (3, c), (4, d)].into_iter().collect();
697        let mates: HashSet<QueueItemId> = [b, d].into_iter().collect();
698
699        bump_to_front(&mut pending, &mates);
700
701        let order: Vec<QueueItemId> = pending.iter().map(|(_, id)| *id).collect();
702        assert_eq!(order, vec![b, d, a, c]);
703    }
704
705    #[test]
706    fn the_track_under_the_cursor_goes_first_and_goes_alone() {
707        let (a, b, c) = (qid(), qid(), qid());
708        let mut q = Queue::default();
709        q.pending.extend([(1, a), (2, b), (3, c)]);
710
711        // Pressed play on the third: it jumps the queue.
712        assert_eq!(next_item(&mut q, Some((3, c))), Some((3, c)));
713        q.in_flight.insert(3, HashSet::from([c]));
714
715        // While it downloads, nothing else starts.
716        assert_eq!(next_item(&mut q, Some((3, c))), None);
717        assert_eq!(q.pending.len(), 2, "the rest wait their turn");
718
719        // Once it has landed the queue runs in order again.
720        q.in_flight.remove(&3);
721        assert_eq!(next_item(&mut q, None), Some((1, a)));
722        assert_eq!(next_item(&mut q, None), Some((2, b)));
723    }
724
725    #[test]
726    fn a_cursor_with_nothing_queued_for_it_holds_nothing_up() {
727        let (a, elsewhere) = (qid(), qid());
728        let mut q = Queue::default();
729        q.pending.push_back((1, a));
730        assert_eq!(next_item(&mut q, Some((9, elsewhere))), Some((1, a)));
731    }
732
733    #[test]
734    fn a_cursor_set_before_its_tracks_were_queued_still_goes_first() {
735        crate::config::isolate_config_for_tests();
736        let item = |title: &str, db_id: i64| crate::player::state::PlaylistItem {
737            playlist_entry_id: None,
738            id: qid(),
739            db_id: Some(db_id),
740            path: std::path::PathBuf::from(format!("/cache/{title}.flac")),
741            title: title.into(),
742            artist: "Artist".into(),
743            album_artist: "Artist".into(),
744            album: "Album".into(),
745            year: None,
746            codec: None,
747            track_number: None,
748            disc: None,
749            duration_ms: None,
750            state: crate::player::state::ItemState::Pending,
751        };
752        let items = vec![item("one", 1), item("two", 2), item("three", 3)];
753        let ids: Vec<_> = items.iter().map(|i| (i.db_id.unwrap(), i.id)).collect();
754
755        // The player has the new queue and its cursor before the download
756        // queue hears of any of it — the first play after launch.
757        let state = SharedPlayerState::new();
758        state.add_items(items);
759        state.set_cursor(Some(ids[0].1));
760
761        let (cmd_tx, _cmd_rx) = crossbeam_channel::unbounded();
762        let inner = Arc::new(Inner {
763            queue: Mutex::new(Queue::default()),
764            has_work: Condvar::new(),
765            state,
766            cmd_tx,
767            log_buf: Arc::new(StdMutex::new(Vec::new())),
768            last_evicted: Mutex::new(None),
769            spawned: std::sync::atomic::AtomicUsize::new(0),
770        });
771        inner.queue.lock().pending.extend(ids.iter().copied());
772        promote_cursor(&inner, ids[0].1);
773
774        let left: Vec<_> = inner.queue.lock().pending.iter().copied().collect();
775        assert_eq!(
776            left,
777            vec![ids[2]],
778            "the cursor's track and the next went to the priority lane"
779        );
780    }
781}