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