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.
465///
466/// Looks before it first waits. The watcher is made with the queue, by the
467/// first downloads queued — on the first play after launch, by that play — and
468/// the player can move the cursor between `enqueue` looking at it and this
469/// thread starting to listen. Waiting first, that move was missed by both, and
470/// the track waited its turn behind whatever else was queued.
471fn cursor_watcher(inner: Arc<Inner>) {
472    let changed = crate::signal::engine_changed();
473    let mut seen = changed.generation();
474    let mut last_cursor: Option<QueueItemId> = None;
475    loop {
476        let current = inner.state.cursor();
477        if current != last_cursor {
478            last_cursor = current;
479            inner.has_work.notify_all();
480            if let Some(cursor_id) = current {
481                promote_cursor(&inner, cursor_id);
482            }
483        }
484        seen = changed.wait(seen);
485    }
486}
487
488/// Send the cursor's track, and the one after it, down the priority lane, if
489/// the cursor's track is waiting in the queue.
490///
491/// Called when the cursor moves and when tracks are queued, because either can
492/// happen first: playing an album sends the new queue to the player and queues
493/// its downloads at once, and the player may set the cursor before or after.
494/// Waiting for the cursor alone missed the first play after launch, when the
495/// queue (and this watcher) did not exist until those downloads made it.
496fn promote_cursor(inner: &Arc<Inner>, cursor_id: QueueItemId) {
497    let is_pending = inner
498        .state
499        .item_load_state(cursor_id)
500        .is_some_and(|s| matches!(s, LoadState::Pending));
501    if !is_pending {
502        return;
503    }
504
505    let album_mate_ids: HashSet<QueueItemId> = inner
506        .state
507        .same_album_item_ids(cursor_id)
508        .into_iter()
509        .collect();
510
511    let mut priority_items = Vec::new();
512    {
513        let mut q = inner.queue.lock();
514        if let Some(pos) = q.pending.iter().position(|(_, qid)| *qid == cursor_id) {
515            priority_items.push(q.pending.remove(pos).expect("position just found"));
516
517            if !album_mate_ids.is_empty() {
518                bump_to_front(&mut q.pending, &album_mate_ids);
519            }
520
521            // Grab the next track too, for gapless lookahead.
522            if let Some(next) = q.pending.pop_front() {
523                priority_items.push(next);
524            }
525        }
526    }
527
528    for item in priority_items {
529        dispatch_priority(inner, item);
530    }
531}
532
533/// The process's download queue.
534///
535/// One player means one pool, one priority lane and one cursor watcher; a
536/// second set would compete with the first for the same link and the same
537/// cursor. Every front end reaches downloads through here — the TUI directly,
538/// the FFI and the GraphQL server through `helpers::spawn_downloads`.
539///
540/// `log_buf` is only honoured by whoever initialises it, which is the TUI when
541/// it is running, since it is the only front end that shows the buffer.
542pub fn shared(
543    cmd_tx: &crossbeam_channel::Sender<PlayerCommand>,
544    state: &Arc<SharedPlayerState>,
545    log_buf: Option<Arc<StdMutex<Vec<String>>>>,
546) -> &'static DownloadQueue {
547    static QUEUE: std::sync::OnceLock<DownloadQueue> = std::sync::OnceLock::new();
548    QUEUE.get_or_init(|| {
549        DownloadQueue::spawn(
550            cmd_tx.clone(),
551            state.clone(),
552            log_buf.unwrap_or_else(|| Arc::new(StdMutex::new(Vec::new()))),
553        )
554    })
555}
556
557#[cfg(test)]
558mod tests {
559    use super::*;
560
561    fn qid() -> QueueItemId {
562        QueueItemId::new()
563    }
564
565    #[test]
566    fn priority_lane_never_exceeds_its_permits() {
567        let mut q = Queue::default();
568
569        // Rapid cursor movement: a fresh track lands on the lane every poll.
570        let mut spawned = 0;
571        for i in 0..500 {
572            if claim_priority(&mut q, (i, qid())) == Dispatch::Spawn {
573                spawned += 1;
574            }
575            assert!(
576                q.priority_active <= PRIORITY_PERMITS,
577                "priority lane over its permit count at iteration {}",
578                i
579            );
580        }
581
582        assert_eq!(spawned, PRIORITY_PERMITS, "only permitted claims may spawn");
583        assert_eq!(
584            q.pending.len(),
585            500 - PRIORITY_PERMITS,
586            "everything else must be queued, not dropped"
587        );
588    }
589
590    #[test]
591    fn released_permits_are_reusable() {
592        let mut q = Queue::default();
593        assert_eq!(claim_priority(&mut q, (1, qid())), Dispatch::Spawn);
594        assert_eq!(claim_priority(&mut q, (2, qid())), Dispatch::Spawn);
595        assert_eq!(claim_priority(&mut q, (3, qid())), Dispatch::Requeued);
596
597        release_priority(&mut q, 1);
598        assert_eq!(claim_priority(&mut q, (4, qid())), Dispatch::Spawn);
599        assert!(q.priority_active <= PRIORITY_PERMITS);
600    }
601
602    #[test]
603    fn an_in_flight_track_is_never_claimed_twice() {
604        let mut q = Queue::default();
605        let id = qid();
606        assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::Spawn);
607        assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::AlreadyRunning);
608        assert_eq!(q.priority_active, 1);
609        assert!(
610            q.pending.is_empty(),
611            "a duplicate request must not re-queue the track"
612        );
613    }
614
615    #[test]
616    fn playing_a_track_again_joins_the_transfer_already_running() {
617        // Playing something twice before it has arrived makes a second queue
618        // entry with an id of its own. The track is the same, and so is the
619        // file a download would write — two of them would truncate and write
620        // over one another, and whichever finished first would rename it away
621        // from the other.
622        let mut q = Queue::default();
623        let (first, again) = (qid(), qid());
624        assert_eq!(claim_priority(&mut q, (7, first)), Dispatch::Spawn);
625        assert_eq!(claim_priority(&mut q, (7, again)), Dispatch::AlreadyRunning);
626
627        assert_eq!(q.priority_active, 1, "one transfer, not two");
628        assert!(q.pending.is_empty());
629        assert_eq!(
630            q.in_flight.get(&7),
631            Some(&HashSet::from([first, again])),
632            "both entries wait on the one transfer"
633        );
634    }
635
636    #[test]
637    fn a_worker_picking_up_a_duplicate_waits_on_the_running_one() {
638        // The same, arriving through the queue rather than the priority lane.
639        let mut q = Queue::default();
640        let (running, queued) = (qid(), qid());
641        assert_eq!(claim_priority(&mut q, (7, running)), Dispatch::Spawn);
642
643        // What `worker_loop` does with the next pending item.
644        match q.in_flight.get_mut(&7) {
645            Some(waiting) => {
646                waiting.insert(queued);
647            }
648            None => panic!("the track should already be claimed"),
649        }
650
651        assert_eq!(
652            q.in_flight.get(&7),
653            Some(&HashSet::from([running, queued])),
654            "the queued entry waits rather than starting a second transfer"
655        );
656    }
657
658    #[test]
659    fn different_tracks_still_run_side_by_side() {
660        // Keying on the track must not serialise unrelated downloads.
661        let mut q = Queue::default();
662        assert_eq!(claim_priority(&mut q, (1, qid())), Dispatch::Spawn);
663        assert_eq!(claim_priority(&mut q, (2, qid())), Dispatch::Spawn);
664        assert_eq!(q.priority_active, 2);
665    }
666
667    #[test]
668    fn requeued_priority_item_goes_to_the_head_of_the_queue() {
669        let mut q = Queue::default();
670        q.pending.push_back((9, qid()));
671        for i in 0..PRIORITY_PERMITS {
672            claim_priority(&mut q, (i as i64, qid()));
673        }
674
675        let wanted = qid();
676        assert_eq!(claim_priority(&mut q, (7, wanted)), Dispatch::Requeued);
677        assert_eq!(q.pending.front().map(|(_, id)| *id), Some(wanted));
678    }
679
680    #[test]
681    fn claiming_removes_a_duplicate_queue_entry() {
682        let mut q = Queue::default();
683        let id = qid();
684        q.pending.push_back((1, id));
685        q.pending.push_back((2, qid()));
686
687        assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::Spawn);
688        assert_eq!(
689            q.pending.len(),
690            1,
691            "the pool must not also pick up the claimed track"
692        );
693    }
694
695    #[test]
696    fn bump_to_front_preserves_relative_order() {
697        let (a, b, c, d) = (qid(), qid(), qid(), qid());
698        let mut pending: VecDeque<(i64, QueueItemId)> =
699            [(1, a), (2, b), (3, c), (4, d)].into_iter().collect();
700        let mates: HashSet<QueueItemId> = [b, d].into_iter().collect();
701
702        bump_to_front(&mut pending, &mates);
703
704        let order: Vec<QueueItemId> = pending.iter().map(|(_, id)| *id).collect();
705        assert_eq!(order, vec![b, d, a, c]);
706    }
707
708    #[test]
709    fn the_track_under_the_cursor_goes_first_and_goes_alone() {
710        let (a, b, c) = (qid(), qid(), qid());
711        let mut q = Queue::default();
712        q.pending.extend([(1, a), (2, b), (3, c)]);
713
714        // Pressed play on the third: it jumps the queue.
715        assert_eq!(next_item(&mut q, Some((3, c))), Some((3, c)));
716        q.in_flight.insert(3, HashSet::from([c]));
717
718        // While it downloads, nothing else starts.
719        assert_eq!(next_item(&mut q, Some((3, c))), None);
720        assert_eq!(q.pending.len(), 2, "the rest wait their turn");
721
722        // Once it has landed the queue runs in order again.
723        q.in_flight.remove(&3);
724        assert_eq!(next_item(&mut q, None), Some((1, a)));
725        assert_eq!(next_item(&mut q, None), Some((2, b)));
726    }
727
728    #[test]
729    fn a_cursor_with_nothing_queued_for_it_holds_nothing_up() {
730        let (a, elsewhere) = (qid(), qid());
731        let mut q = Queue::default();
732        q.pending.push_back((1, a));
733        assert_eq!(next_item(&mut q, Some((9, elsewhere))), Some((1, a)));
734    }
735
736    /// The first play after launch: the player has the new queue and its
737    /// cursor, and the download queue holds its tracks.
738    fn first_play() -> (Arc<Inner>, Vec<(i64, QueueItemId)>) {
739        crate::config::isolate_config_for_tests();
740        let item = |title: &str, db_id: i64| crate::player::state::PlaylistItem {
741            playlist_entry_id: None,
742            id: qid(),
743            db_id: Some(db_id),
744            path: std::path::PathBuf::from(format!("/cache/{title}.flac")),
745            title: title.into(),
746            artist: "Artist".into(),
747            album_artist: "Artist".into(),
748            album: "Album".into(),
749            year: None,
750            codec: None,
751            track_number: None,
752            disc: None,
753            duration_ms: None,
754            state: crate::player::state::ItemState::Pending,
755        };
756        let items = vec![item("one", 1), item("two", 2), item("three", 3)];
757        let ids: Vec<_> = items.iter().map(|i| (i.db_id.unwrap(), i.id)).collect();
758
759        // The player has the new queue and its cursor before the download
760        // queue hears of any of it — the first play after launch.
761        let state = SharedPlayerState::new();
762        state.add_items(items);
763        state.set_cursor(Some(ids[0].1));
764
765        let (cmd_tx, _cmd_rx) = crossbeam_channel::unbounded();
766        let inner = Arc::new(Inner {
767            queue: Mutex::new(Queue::default()),
768            has_work: Condvar::new(),
769            state,
770            cmd_tx,
771            log_buf: Arc::new(StdMutex::new(Vec::new())),
772            last_evicted: Mutex::new(None),
773            spawned: std::sync::atomic::AtomicUsize::new(0),
774        });
775        inner.queue.lock().pending.extend(ids.iter().copied());
776        (inner, ids)
777    }
778
779    #[test]
780    fn a_cursor_set_before_its_tracks_were_queued_still_goes_first() {
781        let (inner, ids) = first_play();
782        promote_cursor(&inner, ids[0].1);
783
784        let left: Vec<_> = inner.queue.lock().pending.iter().copied().collect();
785        assert_eq!(
786            left,
787            vec![ids[2]],
788            "the cursor's track and the next went to the priority lane"
789        );
790    }
791
792    /// The watcher starts after the cursor moved, with nothing further to wake
793    /// it: it still sends the cursor's track ahead.
794    #[test]
795    fn a_watcher_started_after_the_cursor_moved_still_promotes_it() {
796        let (inner, ids) = first_play();
797        let watched = inner.clone();
798        std::thread::spawn(move || cursor_watcher(watched));
799
800        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
801        while inner.queue.lock().pending.len() > 1 && std::time::Instant::now() < deadline {
802            std::thread::sleep(std::time::Duration::from_millis(10));
803        }
804        let left: Vec<_> = inner.queue.lock().pending.iter().copied().collect();
805        assert_eq!(left, vec![ids[2]], "promoted without waiting for a change");
806    }
807}