Skip to main content

koan_core/remote/
queue.rs

1use std::collections::{HashSet, VecDeque};
2use std::panic::AssertUnwindSafe;
3use std::sync::Arc;
4use std::time::Duration;
5
6use parking_lot::{Condvar, Mutex};
7
8use crate::config;
9use crate::db::queries;
10use crate::helpers::download_track;
11use crate::player::commands::PlayerCommand;
12use crate::player::state::{ItemState, QueueItemId, SharedPlayerState};
13use crate::remote::downloads::{self, DownloadStore};
14
15/// Concurrent downloads the priority lane may run outside the worker pool.
16/// Small on purpose: its job is to get the track under the cursor playing, and
17/// every extra request competes with it for the same link.
18const PRIORITY_PERMITS: usize = 2;
19
20/// How often a running transfer asks whether it is still wanted.
21const CANCEL_CHECK: Duration = Duration::from_millis(250);
22
23/// Persistent download queue — lives for the app's lifetime.
24///
25/// Follows the playlist rather than being told about it: whenever the set of
26/// entries waiting for a file changes, each is queued or joins the transfer
27/// already running for its track, and anything no longer in the playlist is
28/// let go. A front end adds tracks to the player and nothing else, so there is
29/// no second request to race the first and no queue of entries the player has
30/// already discarded.
31///
32/// What is in flight, and who waits on it, is the player's download store.
33/// This decides only what is fetched when.
34///
35/// Downloads run on a pool of worker threads. Cursor changes reorder the queue
36/// so the current track downloads first, followed by same-album tracks for
37/// gapless playback; those jump the queue through a permit-limited priority
38/// lane rather than by spawning unbounded threads.
39#[derive(Clone)]
40pub struct DownloadQueue {
41    inner: Arc<Inner>,
42}
43
44struct Inner {
45    /// Held across every change to which entries wait on which transfer —
46    /// joining, letting go, settling — so none can land between another's
47    /// look and its act. May be held while the player's state or its download
48    /// store is locked; neither is ever held while taking this.
49    queue: Mutex<Queue>,
50    has_work: Condvar,
51    state: Arc<SharedPlayerState>,
52    cmd_tx: crossbeam_channel::Sender<PlayerCommand>,
53    /// When the cache was last trimmed to its limit.
54    last_evicted: Mutex<Option<std::time::Instant>>,
55    /// Worker threads started so far. Grows to the configured count as work
56    /// arrives; a worker past the current count parks rather than exits.
57    spawned: std::sync::atomic::AtomicUsize,
58}
59
60/// The parallel-downloads setting as it stands now, not as it stood when the
61/// queue was made: the queue lives as long as the process.
62fn workers_allowed() -> usize {
63    config::Config::cached()
64        .remote
65        .download_workers
66        .clamp(1, 16)
67}
68
69/// Start workers until there are as many as the setting allows.
70fn ensure_workers(inner: &Arc<Inner>) {
71    use std::sync::atomic::Ordering;
72    let want = workers_allowed();
73    loop {
74        let have = inner.spawned.load(Ordering::Relaxed);
75        if have >= want {
76            return;
77        }
78        if inner
79            .spawned
80            .compare_exchange(have, have + 1, Ordering::Relaxed, Ordering::Relaxed)
81            .is_err()
82        {
83            continue;
84        }
85        let worker = inner.clone();
86        if let Err(e) = std::thread::Builder::new()
87            .name(format!("koan-dl-{have}"))
88            .spawn(move || {
89                // The cache is trimmed once at the start, by the first worker
90                // rather than when the queue is made: a player that never
91                // fetches anything has no business reading the cache's size.
92                if have == 0 {
93                    trim_cache(&worker, false);
94                }
95                worker_loop(worker, have)
96            })
97        {
98            log::error!("failed to spawn download worker {have}: {e}");
99            inner.spawned.fetch_sub(1, Ordering::Relaxed);
100            return;
101        }
102    }
103}
104
105/// What a worker should fetch next: the track under the cursor ahead of
106/// everything, nothing at all while that track is already being fetched, and
107/// otherwise the front of the queue.
108///
109/// Tracks wanted only in the cache come after everything in the playlist.
110fn next_item(
111    q: &mut Queue,
112    store: &DownloadStore,
113    cursor: Option<(i64, QueueItemId)>,
114) -> Option<Job> {
115    let entry = match cursor {
116        Some((db_id, _)) if store.in_flight(db_id) => return None,
117        Some((_, queue_id)) => q
118            .pending
119            .iter()
120            .position(|(_, qid)| *qid == queue_id)
121            .and_then(|ix| q.pending.remove(ix))
122            .or_else(|| q.pending.pop_front()),
123        None => q.pending.pop_front(),
124    };
125    entry
126        .map(|(db_id, id)| (db_id, Some(id)))
127        .or_else(|| q.cache.pop_front().map(|db_id| (db_id, None)))
128}
129
130/// Someone asked for music: try the server now rather than when the outage
131/// backoff next says to. The backoff is for retries nobody is waiting on; a
132/// minute of it, earned while the phone was in a pocket with no signal, is not
133/// how long a tap on play should take.
134fn retry_server_now() {
135    if let Some(client) = crate::helpers::subsonic_client(&config::Config::cached()) {
136        client.outage().retry_now();
137    }
138}
139
140/// The track under the cursor, when it still has to be fetched.
141fn cursor_download(state: &SharedPlayerState) -> Option<(i64, QueueItemId)> {
142    let id = state.cursor()?;
143    let item = state.get_item(id)?;
144    if item.state != ItemState::Pending {
145        return None;
146    }
147    Some((item.db_id?, id))
148}
149
150/// A track to fetch, and the queue entry it is for — `None` for a track
151/// wanted only in the cache.
152type Job = (i64, Option<QueueItemId>);
153
154/// What is to be fetched, in the order it will be. What is being fetched is
155/// the download store's.
156#[derive(Default)]
157struct Queue {
158    /// Queue entries whose files are still to be fetched.
159    pending: VecDeque<(i64, QueueItemId)>,
160    /// Tracks wanted in the cache with no queue entry behind them.
161    cache: VecDeque<i64>,
162    priority_active: usize,
163    /// The tracks in the playback window as last worked out, which eviction
164    /// leaves alone. `None` when the whole playlist is wanted: no cache limit,
165    /// or a window that could not be worked out.
166    window: Option<HashSet<i64>>,
167}
168
169/// What a priority request should do, given the state of the lane.
170#[derive(Debug, PartialEq, Eq)]
171enum Dispatch {
172    /// A permit was taken and the item claimed — spawn a thread for it.
173    Spawn,
174    /// No permit free; the item now sits at the head of the work queue.
175    Requeued,
176    /// Some thread is already downloading it.
177    AlreadyRunning,
178}
179
180/// Claim `item` for the priority lane, or push it to the front of the queue if
181/// every permit is taken. On `Spawn` the caller owns a permit and must hand it
182/// back, through a `Permit`, when the download ends.
183fn claim_priority(q: &mut Queue, store: &DownloadStore, item: (i64, QueueItemId)) -> Dispatch {
184    let (db_id, queue_id) = item;
185    q.pending.retain(|(_, qid)| *qid != queue_id);
186
187    // Already being fetched: wait on it rather than fetch it again. The entry
188    // is remembered so it gets the answer when the one transfer lands.
189    if store.join(db_id, Some(queue_id)) {
190        return Dispatch::AlreadyRunning;
191    }
192    if q.priority_active >= PRIORITY_PERMITS {
193        q.pending.push_front(item);
194        return Dispatch::Requeued;
195    }
196    q.priority_active += 1;
197    store.claim(db_id, Some(queue_id));
198    Dispatch::Spawn
199}
200
201/// Hands back a priority permit however the download ends — including a panic.
202struct Permit {
203    inner: Arc<Inner>,
204}
205
206impl Drop for Permit {
207    fn drop(&mut self) {
208        let mut q = self.inner.queue.lock();
209        q.priority_active = q.priority_active.saturating_sub(1);
210        drop(q);
211        self.inner.has_work.notify_all();
212    }
213}
214
215impl DownloadQueue {
216    /// Spawn the download queue and the thread that keeps it following the
217    /// playlist. Workers start with the first thing to fetch.
218    pub fn spawn(
219        cmd_tx: crossbeam_channel::Sender<PlayerCommand>,
220        state: Arc<SharedPlayerState>,
221    ) -> Self {
222        let inner = Arc::new(Inner {
223            queue: Mutex::new(Queue::default()),
224            has_work: Condvar::new(),
225            state,
226            cmd_tx,
227            last_evicted: Mutex::new(None),
228            spawned: std::sync::atomic::AtomicUsize::new(0),
229        });
230
231        let watcher_inner = inner.clone();
232        if let Err(e) = std::thread::Builder::new()
233            .name("koan-dl-watch".into())
234            .spawn(move || follow_playlist(watcher_inner))
235        {
236            log::error!("failed to spawn download playlist watcher: {}", e);
237        }
238
239        Self { inner }
240    }
241
242    /// Fetch these tracks into the cache, with no queue entry to play them.
243    /// They go behind everything the playlist is waiting on.
244    pub fn cache(&self, track_ids: Vec<i64>) {
245        if track_ids.is_empty() {
246            return;
247        }
248        self.inner.queue.lock().cache.extend(track_ids);
249        wake_workers(&self.inner);
250    }
251}
252
253/// There is work: start what workers the setting allows, try the server now,
254/// and wake them.
255fn wake_workers(inner: &Arc<Inner>) {
256    ensure_workers(inner);
257    retry_server_now();
258    inner.has_work.notify_all();
259}
260
261/// Bring the queue in line with the playlist.
262///
263/// Every entry still waiting for its file is either waiting on the transfer
264/// for its track or queued, in the order the player will reach it: from the
265/// cursor to the end, then from the top. Anything the playlist no longer
266/// holds is let go.
267///
268/// With a cache limit, only the entries in the playback window are wanted.
269/// One past it is let go like one removed, and a later sync brings it in when
270/// the cursor or the cache makes room. One inside it is wanted on every sync
271/// that finds it there, so its transfer is never abandoned by a window
272/// recomputed around it.
273///
274/// The playlist is read under the queue lock, as settling writes it: read
275/// before, it could show an entry still pending whose transfer settled a
276/// moment later, and the entry would be fetched a second time. The window is
277/// worked out before, from the database; an entry added in between is
278/// outside it, and the sync its arrival causes brings it in.
279fn sync(inner: &Arc<Inner>) {
280    let window = playback_window(&inner.state);
281    let mut q = inner.queue.lock();
282    let wanted = inner.state.pending_downloads();
283    let entries = window
284        .as_ref()
285        .map(|w| w.iter().map(|(_, id)| *id).collect::<HashSet<_>>());
286    let added = sync_with(&mut q, inner.state.downloads(), &wanted, entries.as_ref());
287    q.window = window.map(|w| w.into_iter().map(|(track, _)| track).collect());
288    drop(q);
289    if added {
290        wake_workers(inner);
291    }
292}
293
294/// The entries the cache has room for, from the cursor on
295/// (`helpers::playback_window`). `None` with no cache limit, and when the
296/// database cannot say: then the whole playlist is fetched, as without one.
297fn playback_window(state: &SharedPlayerState) -> Option<Vec<(i64, QueueItemId)>> {
298    let limit = config::Config::cached().cache_limit_bytes()?;
299    let mut order = state.playback_order();
300    let tracks: Vec<i64> = order.iter().map(|(track, _)| *track).collect();
301    let fits = crate::db::pool::shared()
302        .get()
303        .map_err(|e| e.to_string())
304        .and_then(|db| {
305            crate::helpers::playback_window(&db, limit, &tracks).map_err(|e| e.to_string())
306        });
307    match fits {
308        Ok(fits) => {
309            order.truncate(fits);
310            Some(order)
311        }
312        Err(e) => {
313            log::warn!("cache window: {e}");
314            None
315        }
316    }
317}
318
319/// [`sync`] against a list already read, in the order to fetch it, cut to
320/// `window` when there is one. Whether anything was queued that was not
321/// before.
322fn sync_with(
323    q: &mut Queue,
324    store: &DownloadStore,
325    wanted: &[(i64, QueueItemId)],
326    window: Option<&HashSet<QueueItemId>>,
327) -> bool {
328    let wanted: Vec<_> = match window {
329        Some(window) => wanted
330            .iter()
331            .filter(|(_, id)| window.contains(id))
332            .copied()
333            .collect(),
334        None => wanted.to_vec(),
335    };
336    let unfetched = store.resync(&wanted);
337    let before: HashSet<QueueItemId> = q.pending.iter().map(|(_, id)| *id).collect();
338    let added = unfetched.iter().any(|(_, id)| !before.contains(id));
339    q.pending = unfetched.into();
340    added
341}
342
343/// Start a priority download, or queue it at the front when the lane is full.
344fn dispatch_priority(inner: &Arc<Inner>, item: (i64, QueueItemId)) {
345    let dispatch = claim_priority(&mut inner.queue.lock(), inner.state.downloads(), item);
346    match dispatch {
347        Dispatch::AlreadyRunning => {}
348        Dispatch::Requeued => {
349            inner.has_work.notify_one();
350        }
351        Dispatch::Spawn => {
352            let spawn_inner = inner.clone();
353            let spawned = std::thread::Builder::new()
354                .name("koan-dl-prio".into())
355                .spawn(move || {
356                    let _permit = Permit {
357                        inner: spawn_inner.clone(),
358                    };
359                    run_download(&spawn_inner, item.0);
360                });
361            if let Err(e) = spawned {
362                log::error!("failed to spawn priority download: {}", e);
363                let mut q = inner.queue.lock();
364                let store = inner.state.downloads();
365                for id in downloads::withdraw(store, item.0) {
366                    q.pending.push_front((item.0, id));
367                }
368                q.priority_active = q.priority_active.saturating_sub(1);
369                drop(q);
370                inner.has_work.notify_one();
371            }
372        }
373    }
374}
375
376/// Run one download, containing any panic so the worker pool never shrinks,
377/// and settle every entry waiting on it. The transfer was claimed in the
378/// player's store before this was called.
379///
380/// The client is looked up per download, not once for the queue's lifetime:
381/// the queue lives as long as the process, and signing in, out or elsewhere
382/// has to reach it.
383fn run_download(inner: &Arc<Inner>, db_id: i64) {
384    let store = inner.state.downloads();
385    let cfg = config::Config::cached();
386    let result = match crate::helpers::subsonic_client(&cfg) {
387        // Failed, not left Pending: the player waits for Ready, so a queue of
388        // tracks that can never arrive would otherwise sit saying nothing.
389        None => Some(Err(crate::helpers::remote_unavailable(&cfg))),
390        Some(client) => {
391            // Asked per chunk, answered from the store at most every
392            // `CANCEL_CHECK`. Replacing the queue is one command and one
393            // playlist change (`SharedPlayerState::replace_playlist`), so the
394            // playlist is never momentarily without a track still wanted.
395            let checked = std::cell::Cell::new(std::time::Instant::now());
396            let cancelled = || {
397                if checked.get().elapsed() < CANCEL_CHECK {
398                    return false;
399                }
400                checked.set(std::time::Instant::now());
401                store.abandoned(db_id)
402            };
403            std::panic::catch_unwind(AssertUnwindSafe(|| {
404                download_track(
405                    db_id,
406                    &cancelled,
407                    &inner.cmd_tx,
408                    &inner.state,
409                    &cfg,
410                    &client,
411                )
412            }))
413            .unwrap_or_else(|_| {
414                log::error!("download panicked for track {db_id}");
415                Some(Err("download panicked".into()))
416            })
417        }
418    };
419
420    // Under the queue lock, so no entry can join the transfer between being
421    // told and the transfer ending: one arriving after this finds no
422    // transfer and is queued afresh.
423    if let Some(Ok(_)) = &result
424        && store.kept(db_id)
425    {
426        pin(db_id);
427    }
428    let mut q = inner.queue.lock();
429    let settled = match result {
430        Some(result) => Some(downloads::settle(&inner.state, db_id, &result)),
431        // Withdrawn because nothing wanted it, as far as it last looked. An
432        // entry that has asked since goes back to the front.
433        None => {
434            for id in downloads::withdraw(store, db_id) {
435                q.pending.push_front((db_id, id));
436            }
437            None
438        }
439    };
440    drop(q);
441    if let Some(settled) = settled {
442        settled.announce(&inner.cmd_tx);
443    }
444    // Workers held back for the track under the cursor go on from here.
445    inner.has_work.notify_all();
446    // The file's real size replaces its estimate, which may move the window.
447    if config::Config::cached().cache_limit_bytes().is_some() {
448        sync(inner);
449    }
450    trim_cache(inner, false);
451}
452
453/// A track downloaded because someone asked for the file: eviction takes it
454/// only once everything fetched for playback has gone.
455fn pin(db_id: i64) {
456    let pinned = crate::db::pool::shared()
457        .get()
458        .map_err(|e| e.to_string())
459        .and_then(|db| queries::pin_cached(&db.conn, &[db_id]).map_err(|e| e.to_string()));
460    if let Err(e) = pinned {
461        log::warn!("could not pin the download of track {db_id}: {e}");
462    }
463}
464
465/// How often the cache is checked against its limit, at most: each download
466/// adds to it, and a check reads the whole cache's size from the database.
467const EVICT_EVERY: std::time::Duration = std::time::Duration::from_secs(60);
468
469/// Trim the cache to its configured limit, keeping the playback window: the
470/// player may be reading those files, and they are what is wanted next.
471/// Without a window, everything in the playlist is kept. `now` skips the
472/// throttle, for a limit just changed. The limit is read afresh, so one set
473/// in Settings applies without a restart.
474fn trim_cache(inner: &Inner, now: bool) {
475    {
476        let mut last = inner.last_evicted.lock();
477        if !now && last.is_some_and(|t| t.elapsed() < EVICT_EVERY) {
478            return;
479        }
480        *last = Some(std::time::Instant::now());
481    }
482    let cfg = config::Config::cached();
483    if cfg.cache_limit_bytes().is_none() {
484        return;
485    }
486    let keep = inner.queue.lock().window.clone().unwrap_or_else(|| {
487        inner
488            .state
489            .playback_order()
490            .into_iter()
491            .map(|(track, _)| track)
492            .collect()
493    });
494    match crate::db::pool::shared().get() {
495        Ok(db) => {
496            if crate::helpers::evict_cache(&db, &cfg, &keep, false) > 0 {
497                // Played entries pointed at the files just removed. Pending
498                // again, they are fetched when the window comes back to them.
499                inner.state.reset_items_with_missing_files();
500            }
501        }
502        Err(e) => log::warn!("cache eviction: could not open the database: {e}"),
503    }
504}
505
506/// Moved whenever the cache limit is changed.
507static LIMIT_CHANGES: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
508
509/// The cache limit was changed: work the window out again and evict down to
510/// the new limit now, rather than at the next download.
511pub fn cache_limit_changed() {
512    LIMIT_CHANGES.fetch_add(1, std::sync::atomic::Ordering::AcqRel);
513    crate::signal::engine_changed().bump();
514}
515
516/// Worker loop: wait for work, download, repeat.
517///
518/// While the server is down, workers take nothing new: the download that found
519/// it down waits it out, and everything behind it stays `Pending` rather than
520/// each piling onto a server that is not answering.
521/// One download worker. `index` is its place in the pool: a worker past the
522/// current parallel-downloads setting parks until the setting lets it work.
523///
524/// While the track under the cursor is still being fetched, no worker starts
525/// another transfer. On a slow link every parallel download takes a share of
526/// the bandwidth, and the one track someone is waiting to hear would arrive
527/// last among equals; the rest of the album can follow it.
528fn worker_loop(inner: Arc<Inner>, index: usize) {
529    loop {
530        let item = loop {
531            if let Some(client) = crate::helpers::subsonic_client(&config::Config::cached()) {
532                client.outage().hold();
533            }
534            let cursor = cursor_download(&inner.state);
535            let mut q = inner.queue.lock();
536            if index >= workers_allowed() {
537                inner.has_work.wait(&mut q);
538                continue;
539            }
540            let store = inner.state.downloads();
541            match next_item(&mut q, store, cursor) {
542                // Already being fetched: the entry waits on that transfer
543                // rather than starting a second over it.
544                Some((db_id, entry)) => {
545                    if store.claim(db_id, entry) {
546                        break db_id;
547                    }
548                }
549                None => inner.has_work.wait(&mut q),
550            }
551        };
552        run_download(&inner, item);
553    }
554}
555
556/// Keep the queue following the player: when the set of entries waiting for a
557/// file may have changed, bring the queue in line with it; when the cursor moves to a pending track, hand it
558/// and the next track to the priority lane and bump same-album tracks to the
559/// front.
560///
561/// A cursor move also wakes the workers, which hold back while the cursor's
562/// track is being fetched: one that looked before the cursor moved is waiting
563/// on a track nobody wants any more.
564///
565/// Looks before it first waits, so a playlist the player already holds when
566/// the queue is made — a restored session — is fetched without waiting for
567/// something else to change.
568fn follow_playlist(inner: Arc<Inner>) {
569    let changed = crate::signal::engine_changed();
570    let mut seen = changed.generation();
571    let mut last_version: Option<u64> = None;
572    let mut last_cursor: Option<QueueItemId> = None;
573    let mut last_limit = LIMIT_CHANGES.load(std::sync::atomic::Ordering::Acquire);
574    loop {
575        // The queue runs from the cursor, so a cursor move reorders it as
576        // well. Synced before promoting, so the cursor's track is queued by
577        // the time it is looked for.
578        let version = inner.state.pending_version();
579        let current = inner.state.cursor();
580        let limit = LIMIT_CHANGES.load(std::sync::atomic::Ordering::Acquire);
581        if last_version != Some(version) || current != last_cursor || limit != last_limit {
582            last_version = Some(version);
583            sync(&inner);
584        }
585        if limit != last_limit {
586            last_limit = limit;
587            trim_cache(&inner, true);
588        }
589        if current != last_cursor {
590            last_cursor = current;
591            inner.has_work.notify_all();
592            if let Some(cursor_id) = current {
593                promote_cursor(&inner, cursor_id);
594            }
595        }
596        seen = changed.wait(seen);
597    }
598}
599
600/// Send the cursor's track, and the one after it, down the priority lane, if
601/// the cursor's track is waiting in the queue. The queue is already in the
602/// order the player will reach it, so the one after is its front.
603fn promote_cursor(inner: &Arc<Inner>, cursor_id: QueueItemId) {
604    if inner.state.item_state(cursor_id) != Some(ItemState::Pending) {
605        return;
606    }
607    let mut priority_items = Vec::new();
608    {
609        let mut q = inner.queue.lock();
610        if let Some(pos) = q.pending.iter().position(|(_, qid)| *qid == cursor_id) {
611            priority_items.push(q.pending.remove(pos).expect("position just found"));
612            // The next track too, for gapless lookahead.
613            if let Some(next) = q.pending.pop_front() {
614                priority_items.push(next);
615            }
616        }
617    }
618    for item in priority_items {
619        dispatch_priority(inner, item);
620    }
621}
622
623#[cfg(test)]
624mod tests {
625    use super::*;
626    use crate::player::state::PlaylistItem;
627
628    fn qid() -> QueueItemId {
629        QueueItemId::new()
630    }
631
632    fn waiters(store: &DownloadStore, db_id: i64) -> HashSet<QueueItemId> {
633        store.waiters(db_id).into_iter().collect()
634    }
635
636    #[test]
637    fn priority_lane_never_exceeds_its_permits() {
638        let (mut q, store) = (Queue::default(), DownloadStore::new());
639
640        // Rapid cursor movement: a fresh track lands on the lane every poll.
641        let mut spawned = 0;
642        for i in 0..500 {
643            if claim_priority(&mut q, &store, (i, qid())) == Dispatch::Spawn {
644                spawned += 1;
645            }
646            assert!(
647                q.priority_active <= PRIORITY_PERMITS,
648                "priority lane over its permit count at iteration {}",
649                i
650            );
651        }
652
653        assert_eq!(spawned, PRIORITY_PERMITS, "only permitted claims may spawn");
654        assert_eq!(
655            q.pending.len(),
656            500 - PRIORITY_PERMITS,
657            "everything else must be queued, not dropped"
658        );
659    }
660
661    #[test]
662    fn released_permits_are_reusable() {
663        let (mut q, store) = (Queue::default(), DownloadStore::new());
664        assert_eq!(claim_priority(&mut q, &store, (1, qid())), Dispatch::Spawn);
665        assert_eq!(claim_priority(&mut q, &store, (2, qid())), Dispatch::Spawn);
666        assert_eq!(
667            claim_priority(&mut q, &store, (3, qid())),
668            Dispatch::Requeued
669        );
670
671        let _ = downloads::withdraw(&store, 1);
672        q.priority_active -= 1;
673        assert_eq!(claim_priority(&mut q, &store, (4, qid())), Dispatch::Spawn);
674        assert!(q.priority_active <= PRIORITY_PERMITS);
675    }
676
677    #[test]
678    fn an_in_flight_track_is_never_claimed_twice() {
679        let (mut q, store) = (Queue::default(), DownloadStore::new());
680        let id = qid();
681        assert_eq!(claim_priority(&mut q, &store, (1, id)), Dispatch::Spawn);
682        assert_eq!(
683            claim_priority(&mut q, &store, (1, id)),
684            Dispatch::AlreadyRunning
685        );
686        assert_eq!(q.priority_active, 1);
687        assert!(
688            q.pending.is_empty(),
689            "a duplicate request must not re-queue the track"
690        );
691    }
692
693    #[test]
694    fn playing_a_track_again_joins_the_transfer_already_running() {
695        // Playing something twice before it has arrived makes a second queue
696        // entry with an id of its own. The track is the same, and so is the
697        // file a download would write — two of them would truncate and write
698        // over one another, and whichever finished first would rename it away
699        // from the other.
700        let (mut q, store) = (Queue::default(), DownloadStore::new());
701        let (first, again) = (qid(), qid());
702        assert_eq!(claim_priority(&mut q, &store, (7, first)), Dispatch::Spawn);
703        assert_eq!(
704            claim_priority(&mut q, &store, (7, again)),
705            Dispatch::AlreadyRunning
706        );
707
708        assert_eq!(q.priority_active, 1, "one transfer, not two");
709        assert!(q.pending.is_empty());
710        assert_eq!(
711            waiters(&store, 7),
712            HashSet::from([first, again]),
713            "both entries wait on the one transfer"
714        );
715    }
716
717    #[test]
718    fn different_tracks_still_run_side_by_side() {
719        // Keying on the track must not serialise unrelated downloads.
720        let (mut q, store) = (Queue::default(), DownloadStore::new());
721        assert_eq!(claim_priority(&mut q, &store, (1, qid())), Dispatch::Spawn);
722        assert_eq!(claim_priority(&mut q, &store, (2, qid())), Dispatch::Spawn);
723        assert_eq!(q.priority_active, 2);
724    }
725
726    #[test]
727    fn requeued_priority_item_goes_to_the_head_of_the_queue() {
728        let (mut q, store) = (Queue::default(), DownloadStore::new());
729        q.pending.push_back((9, qid()));
730        for i in 0..PRIORITY_PERMITS {
731            claim_priority(&mut q, &store, (i as i64, qid()));
732        }
733
734        let wanted = qid();
735        assert_eq!(
736            claim_priority(&mut q, &store, (7, wanted)),
737            Dispatch::Requeued
738        );
739        assert_eq!(q.pending.front().map(|(_, id)| *id), Some(wanted));
740    }
741
742    #[test]
743    fn claiming_removes_a_duplicate_queue_entry() {
744        let (mut q, store) = (Queue::default(), DownloadStore::new());
745        let id = qid();
746        q.pending.push_back((1, id));
747        q.pending.push_back((2, qid()));
748
749        assert_eq!(claim_priority(&mut q, &store, (1, id)), Dispatch::Spawn);
750        assert_eq!(
751            q.pending.len(),
752            1,
753            "the pool must not also pick up the claimed track"
754        );
755    }
756
757    #[test]
758    fn the_track_under_the_cursor_goes_first_and_goes_alone() {
759        let (a, b, c) = (qid(), qid(), qid());
760        let (mut q, store) = (Queue::default(), DownloadStore::new());
761        q.pending.extend([(1, a), (2, b), (3, c)]);
762
763        // Pressed play on the third: it jumps the queue.
764        assert_eq!(next_item(&mut q, &store, Some((3, c))), Some((3, Some(c))));
765        store.claim(3, Some(c));
766
767        // While it downloads, nothing else starts.
768        assert_eq!(next_item(&mut q, &store, Some((3, c))), None);
769        assert_eq!(q.pending.len(), 2, "the rest wait their turn");
770
771        // Once it has landed the queue runs in order again.
772        let _ = downloads::withdraw(&store, 3);
773        assert_eq!(next_item(&mut q, &store, None), Some((1, Some(a))));
774        assert_eq!(next_item(&mut q, &store, None), Some((2, Some(b))));
775    }
776
777    #[test]
778    fn a_cursor_with_nothing_queued_for_it_holds_nothing_up() {
779        let (a, elsewhere) = (qid(), qid());
780        let (mut q, store) = (Queue::default(), DownloadStore::new());
781        q.pending.push_back((1, a));
782        assert_eq!(
783            next_item(&mut q, &store, Some((9, elsewhere))),
784            Some((1, Some(a)))
785        );
786    }
787
788    #[test]
789    fn a_track_wanted_only_in_the_cache_waits_behind_the_playlist() {
790        let a = qid();
791        let (mut q, store) = (Queue::default(), DownloadStore::new());
792        q.cache.push_back(5);
793        q.pending.push_back((1, a));
794        assert_eq!(next_item(&mut q, &store, None), Some((1, Some(a))));
795        assert_eq!(next_item(&mut q, &store, None), Some((5, None)));
796    }
797
798    #[test]
799    fn the_queue_lets_go_of_what_the_playlist_no_longer_holds() {
800        // Replacing a large queue used to leave every old entry queued, and
801        // each cost a worker a wait before it gave up.
802        let (old, kept, waiter) = (qid(), qid(), qid());
803        let (mut q, store) = (Queue::default(), DownloadStore::new());
804        q.pending.extend([(1, old), (2, kept)]);
805        store.claim(3, Some(waiter));
806
807        sync_with(&mut q, &store, &[(2, kept)], None);
808
809        assert_eq!(q.pending, VecDeque::from([(2, kept)]));
810        assert!(waiters(&store, 3).is_empty());
811        assert!(
812            store.abandoned(3),
813            "a transfer nothing waits on any more is on its way to being stopped"
814        );
815    }
816
817    #[test]
818    fn the_queue_is_in_the_order_it_is_given() {
819        // The playlist's order from the cursor on, which `pending_downloads`
820        // gives: a queue that kept its old order fetched the tracks before a
821        // cursor that had jumped ahead first.
822        let (a, b, c) = (qid(), qid(), qid());
823        let (mut q, store) = (Queue::default(), DownloadStore::new());
824        q.pending.extend([(1, a), (2, b), (3, c)]);
825
826        assert!(
827            !sync_with(&mut q, &store, &[(3, c), (1, a), (2, b)], None),
828            "nothing new"
829        );
830        assert_eq!(q.pending, VecDeque::from([(3, c), (1, a), (2, b)]));
831
832        let d = qid();
833        assert!(sync_with(
834            &mut q,
835            &store,
836            &[(3, c), (4, d), (1, a), (2, b)],
837            None
838        ));
839        assert_eq!(q.pending, VecDeque::from([(3, c), (4, d), (1, a), (2, b)]));
840    }
841
842    #[test]
843    fn playing_from_the_middle_fetches_from_there_first() {
844        // Pressed play on track three of five: three, four and five come
845        // before one and two.
846        let (inner, ids, _) = queue_over(&[1, 2, 3, 4, 5]);
847        inner.state.set_cursor(Some(ids[2].1));
848        let wanted = inner.state.pending_downloads();
849        sync_with(
850            &mut inner.queue.lock(),
851            inner.state.downloads(),
852            &wanted,
853            None,
854        );
855
856        let order: Vec<i64> = inner.queue.lock().pending.iter().map(|(t, _)| *t).collect();
857        assert_eq!(order, vec![3, 4, 5, 1, 2]);
858    }
859
860    #[test]
861    fn a_new_entry_for_a_track_in_flight_waits_on_that_transfer() {
862        let (running, again) = (qid(), qid());
863        let (mut q, store) = (Queue::default(), DownloadStore::new());
864        store.claim(7, Some(running));
865
866        assert!(!sync_with(
867            &mut q,
868            &store,
869            &[(7, running), (7, again)],
870            None
871        ));
872
873        assert!(q.pending.is_empty());
874        assert_eq!(waiters(&store, 7), HashSet::from([running, again]));
875    }
876
877    #[test]
878    fn a_queue_replaced_with_the_same_track_keeps_its_transfer() {
879        // A new queue holding the same track is new entries for it; the one
880        // transfer serves them and is never wanted by nothing in between.
881        let (old, new) = (qid(), qid());
882        let (mut q, store) = (Queue::default(), DownloadStore::new());
883        store.claim(7, Some(old));
884
885        sync_with(&mut q, &store, &[(7, new)], None);
886
887        assert_eq!(waiters(&store, 7), HashSet::from([new]));
888        assert!(!store.abandoned(7));
889    }
890
891    #[test]
892    fn entries_past_the_window_wait_until_it_reaches_them() {
893        let (played, next, beyond) = (qid(), qid(), qid());
894        let (mut q, store) = (Queue::default(), DownloadStore::new());
895        let wanted = [(2, next), (3, beyond), (1, played)];
896
897        sync_with(&mut q, &store, &wanted, Some(&HashSet::from([next])));
898        assert_eq!(q.pending, VecDeque::from([(2, next)]));
899
900        // The cursor moves on and the window with it.
901        sync_with(&mut q, &store, &wanted, Some(&HashSet::from([beyond])));
902        assert_eq!(q.pending, VecDeque::from([(3, beyond)]));
903    }
904
905    #[test]
906    fn a_track_inside_the_window_keeps_its_transfer_when_the_window_is_recomputed() {
907        let (playing, next, beyond) = (qid(), qid(), qid());
908        let (mut q, store) = (Queue::default(), DownloadStore::new());
909        store.claim(2, Some(next));
910        store.claim(3, Some(beyond));
911        let wanted = [(1, playing), (2, next), (3, beyond)];
912
913        // The download of the next track lands, the window is worked out
914        // again from the cache's new size, and comes out smaller.
915        sync_with(
916            &mut q,
917            &store,
918            &wanted,
919            Some(&HashSet::from([playing, next])),
920        );
921
922        assert!(!store.abandoned(2), "inside the window: still wanted");
923        assert_eq!(waiters(&store, 2), HashSet::from([next]));
924        assert!(store.abandoned(3), "past the window: let go");
925        assert_eq!(q.pending, VecDeque::from([(1, playing)]));
926    }
927
928    /// A player holding one pending item per track id, the cursor on the
929    /// first, and a download queue over it that has heard of none of them.
930    fn queue_over(
931        tracks: &[i64],
932    ) -> (
933        Arc<Inner>,
934        Vec<(i64, QueueItemId)>,
935        crossbeam_channel::Receiver<PlayerCommand>,
936    ) {
937        crate::config::isolate_config_for_tests();
938        let item = |db_id: i64| PlaylistItem {
939            playlist_entry_id: None,
940            id: qid(),
941            db_id: Some(db_id),
942            path: std::path::PathBuf::from(format!("/cache/track-{db_id}.flac")),
943            title: format!("track-{db_id}"),
944            artist: "Artist".into(),
945            album_artist: "Artist".into(),
946            album: "Album".into(),
947            year: None,
948            codec: None,
949            track_number: None,
950            disc: None,
951            duration_ms: None,
952            state: ItemState::Pending,
953            pre_shuffle: None,
954        };
955        let items: Vec<_> = tracks.iter().map(|&db_id| item(db_id)).collect();
956        let ids: Vec<_> = items.iter().map(|i| (i.db_id.unwrap(), i.id)).collect();
957
958        let state = SharedPlayerState::new();
959        state.add_items(items);
960        state.set_cursor(Some(ids[0].1));
961
962        let (cmd_tx, cmd_rx) = crossbeam_channel::unbounded();
963        let inner = Arc::new(Inner {
964            queue: Mutex::new(Queue::default()),
965            has_work: Condvar::new(),
966            state,
967            cmd_tx,
968            last_evicted: Mutex::new(None),
969            spawned: std::sync::atomic::AtomicUsize::new(0),
970        });
971        (inner, ids, cmd_rx)
972    }
973
974    /// The first play after launch: the player has the new queue and its
975    /// cursor, and the download queue holds its tracks.
976    fn first_play() -> (Arc<Inner>, Vec<(i64, QueueItemId)>) {
977        let (inner, ids, _) = queue_over(&[1, 2, 3]);
978        inner.queue.lock().pending.extend(ids.iter().copied());
979        (inner, ids)
980    }
981
982    #[test]
983    fn a_cursor_set_before_its_tracks_were_queued_still_goes_first() {
984        let (inner, ids) = first_play();
985        promote_cursor(&inner, ids[0].1);
986
987        let left: Vec<_> = inner.queue.lock().pending.iter().copied().collect();
988        assert_eq!(
989            left,
990            vec![ids[2]],
991            "the cursor's track and the next went to the priority lane"
992        );
993    }
994
995    fn failed(inner: &Inner, id: QueueItemId) -> bool {
996        matches!(inner.state.item_state(id), Some(ItemState::Failed(_)))
997    }
998
999    /// The watcher starts after the player already has its queue and cursor,
1000    /// with nothing further to wake it: it finds the tracks by itself and
1001    /// fetches every one. With no server, each fails rather than waiting.
1002    #[test]
1003    fn a_watcher_started_over_a_playlist_fetches_it_unprompted() {
1004        let (inner, ids, _) = queue_over(&[1, 2, 3]);
1005        let watched = inner.clone();
1006        std::thread::spawn(move || follow_playlist(watched));
1007
1008        let deadline = std::time::Instant::now() + Duration::from_secs(5);
1009        while !ids.iter().all(|(_, id)| failed(&inner, *id)) && std::time::Instant::now() < deadline
1010        {
1011            std::thread::sleep(Duration::from_millis(10));
1012        }
1013        assert!(ids.iter().all(|(_, id)| failed(&inner, *id)));
1014    }
1015
1016    /// A track queued twice, both entries waiting on its one transfer, the
1017    /// second under the cursor. The transfer cannot happen: no server.
1018    fn duplicate_that_fails() -> (
1019        Arc<Inner>,
1020        crossbeam_channel::Receiver<PlayerCommand>,
1021        QueueItemId,
1022        QueueItemId,
1023    ) {
1024        let (inner, ids, cmd_rx) = queue_over(&[1, 1]);
1025        let (first, again) = (ids[0].1, ids[1].1);
1026        inner.state.set_cursor(Some(again));
1027        let store = inner.state.downloads();
1028        store.claim(1, Some(first));
1029        store.join(1, Some(again));
1030        (inner, cmd_rx, first, again)
1031    }
1032
1033    #[test]
1034    fn a_failed_transfer_fails_every_entry_waiting_on_it() {
1035        let (inner, cmd_rx, first, again) = duplicate_that_fails();
1036
1037        run_download(&inner, 1);
1038
1039        for id in [first, again] {
1040            assert!(failed(&inner, id), "every waiter hears the one answer");
1041        }
1042        let sent: Vec<_> = cmd_rx.try_iter().collect();
1043        assert!(
1044            matches!(sent.as_slice(), [PlayerCommand::TrackFailed(id)] if *id == again),
1045            "the cursor's entry is told it failed, not that it is ready: {sent:?}"
1046        );
1047        assert!(!inner.state.downloads().in_flight(1));
1048    }
1049
1050    #[test]
1051    fn a_transfer_whose_first_entry_was_removed_still_answers_the_rest() {
1052        let (inner, _cmd_rx, first, again) = duplicate_that_fails();
1053        inner.state.remove_item(first);
1054        sync(&inner);
1055
1056        run_download(&inner, 1);
1057
1058        assert!(failed(&inner, again));
1059    }
1060
1061    #[test]
1062    fn a_transfer_withdrawn_with_an_entry_waiting_queues_it_again() {
1063        // The entry joined after the download decided nothing wanted it.
1064        let (inner, ids, _) = queue_over(&[1]);
1065        let store = inner.state.downloads();
1066        store.claim(1, Some(ids[0].1));
1067
1068        let mut q = inner.queue.lock();
1069        for id in downloads::withdraw(store, 1) {
1070            q.pending.push_front((1, id));
1071        }
1072        assert_eq!(q.pending, VecDeque::from([ids[0]]));
1073    }
1074}