use std::collections::{HashSet, VecDeque};
use std::panic::AssertUnwindSafe;
use std::sync::Arc;
use std::time::Duration;
use parking_lot::{Condvar, Mutex};
use crate::config;
use crate::db::queries;
use crate::helpers::download_track;
use crate::player::commands::PlayerCommand;
use crate::player::state::{ItemState, QueueItemId, SharedPlayerState};
use crate::remote::downloads::{self, DownloadStore};
const PRIORITY_PERMITS: usize = 2;
const CANCEL_CHECK: Duration = Duration::from_millis(250);
#[derive(Clone)]
pub struct DownloadQueue {
inner: Arc<Inner>,
}
struct Inner {
queue: Mutex<Queue>,
has_work: Condvar,
state: Arc<SharedPlayerState>,
cmd_tx: crossbeam_channel::Sender<PlayerCommand>,
last_evicted: Mutex<Option<std::time::Instant>>,
spawned: std::sync::atomic::AtomicUsize,
}
fn workers_allowed() -> usize {
config::Config::cached()
.remote
.download_workers
.clamp(1, 16)
}
fn ensure_workers(inner: &Arc<Inner>) {
use std::sync::atomic::Ordering;
let want = workers_allowed();
loop {
let have = inner.spawned.load(Ordering::Relaxed);
if have >= want {
return;
}
if inner
.spawned
.compare_exchange(have, have + 1, Ordering::Relaxed, Ordering::Relaxed)
.is_err()
{
continue;
}
let worker = inner.clone();
if let Err(e) = std::thread::Builder::new()
.name(format!("koan-dl-{have}"))
.spawn(move || {
if have == 0 {
trim_cache(&worker, false);
}
worker_loop(worker, have)
})
{
log::error!("failed to spawn download worker {have}: {e}");
inner.spawned.fetch_sub(1, Ordering::Relaxed);
return;
}
}
}
fn next_item(
q: &mut Queue,
store: &DownloadStore,
cursor: Option<(i64, QueueItemId)>,
) -> Option<Job> {
let entry = match cursor {
Some((db_id, _)) if store.in_flight(db_id) => return None,
Some((_, queue_id)) => q
.pending
.iter()
.position(|(_, qid)| *qid == queue_id)
.and_then(|ix| q.pending.remove(ix))
.or_else(|| q.pending.pop_front()),
None => q.pending.pop_front(),
};
entry
.map(|(db_id, id)| (db_id, Some(id)))
.or_else(|| q.cache.pop_front().map(|db_id| (db_id, None)))
}
fn retry_server_now() {
if let Some(client) = crate::helpers::subsonic_client(&config::Config::cached()) {
client.outage().retry_now();
}
}
fn cursor_download(state: &SharedPlayerState) -> Option<(i64, QueueItemId)> {
let id = state.cursor()?;
let item = state.get_item(id)?;
if item.state != ItemState::Pending {
return None;
}
Some((item.db_id?, id))
}
type Job = (i64, Option<QueueItemId>);
#[derive(Default)]
struct Queue {
pending: VecDeque<(i64, QueueItemId)>,
cache: VecDeque<i64>,
priority_active: usize,
window: Option<HashSet<i64>>,
}
#[derive(Debug, PartialEq, Eq)]
enum Dispatch {
Spawn,
Requeued,
AlreadyRunning,
}
fn claim_priority(q: &mut Queue, store: &DownloadStore, item: (i64, QueueItemId)) -> Dispatch {
let (db_id, queue_id) = item;
q.pending.retain(|(_, qid)| *qid != queue_id);
if store.join(db_id, Some(queue_id)) {
return Dispatch::AlreadyRunning;
}
if q.priority_active >= PRIORITY_PERMITS {
q.pending.push_front(item);
return Dispatch::Requeued;
}
q.priority_active += 1;
store.claim(db_id, Some(queue_id));
Dispatch::Spawn
}
struct Permit {
inner: Arc<Inner>,
}
impl Drop for Permit {
fn drop(&mut self) {
let mut q = self.inner.queue.lock();
q.priority_active = q.priority_active.saturating_sub(1);
drop(q);
self.inner.has_work.notify_all();
}
}
impl DownloadQueue {
pub fn spawn(
cmd_tx: crossbeam_channel::Sender<PlayerCommand>,
state: Arc<SharedPlayerState>,
) -> Self {
let inner = Arc::new(Inner {
queue: Mutex::new(Queue::default()),
has_work: Condvar::new(),
state,
cmd_tx,
last_evicted: Mutex::new(None),
spawned: std::sync::atomic::AtomicUsize::new(0),
});
let watcher_inner = inner.clone();
if let Err(e) = std::thread::Builder::new()
.name("koan-dl-watch".into())
.spawn(move || follow_playlist(watcher_inner))
{
log::error!("failed to spawn download playlist watcher: {}", e);
}
Self { inner }
}
pub fn cache(&self, track_ids: Vec<i64>) {
if track_ids.is_empty() {
return;
}
self.inner.queue.lock().cache.extend(track_ids);
wake_workers(&self.inner);
}
}
fn wake_workers(inner: &Arc<Inner>) {
ensure_workers(inner);
retry_server_now();
inner.has_work.notify_all();
}
fn sync(inner: &Arc<Inner>) {
let window = playback_window(&inner.state);
let mut q = inner.queue.lock();
let wanted = inner.state.pending_downloads();
let entries = window
.as_ref()
.map(|w| w.iter().map(|(_, id)| *id).collect::<HashSet<_>>());
let added = sync_with(&mut q, inner.state.downloads(), &wanted, entries.as_ref());
q.window = window.map(|w| w.into_iter().map(|(track, _)| track).collect());
drop(q);
if added {
wake_workers(inner);
}
}
fn playback_window(state: &SharedPlayerState) -> Option<Vec<(i64, QueueItemId)>> {
let limit = config::Config::cached().cache_limit_bytes()?;
let mut order = state.playback_order();
let tracks: Vec<i64> = order.iter().map(|(track, _)| *track).collect();
let fits = crate::db::pool::shared()
.get()
.map_err(|e| e.to_string())
.and_then(|db| {
crate::helpers::playback_window(&db, limit, &tracks).map_err(|e| e.to_string())
});
match fits {
Ok(fits) => {
order.truncate(fits);
Some(order)
}
Err(e) => {
log::warn!("cache window: {e}");
None
}
}
}
fn sync_with(
q: &mut Queue,
store: &DownloadStore,
wanted: &[(i64, QueueItemId)],
window: Option<&HashSet<QueueItemId>>,
) -> bool {
let wanted: Vec<_> = match window {
Some(window) => wanted
.iter()
.filter(|(_, id)| window.contains(id))
.copied()
.collect(),
None => wanted.to_vec(),
};
let unfetched = store.resync(&wanted);
let before: HashSet<QueueItemId> = q.pending.iter().map(|(_, id)| *id).collect();
let added = unfetched.iter().any(|(_, id)| !before.contains(id));
q.pending = unfetched.into();
added
}
fn dispatch_priority(inner: &Arc<Inner>, item: (i64, QueueItemId)) {
let dispatch = claim_priority(&mut inner.queue.lock(), inner.state.downloads(), item);
match dispatch {
Dispatch::AlreadyRunning => {}
Dispatch::Requeued => {
inner.has_work.notify_one();
}
Dispatch::Spawn => {
let spawn_inner = inner.clone();
let spawned = std::thread::Builder::new()
.name("koan-dl-prio".into())
.spawn(move || {
let _permit = Permit {
inner: spawn_inner.clone(),
};
run_download(&spawn_inner, item.0);
});
if let Err(e) = spawned {
log::error!("failed to spawn priority download: {}", e);
let mut q = inner.queue.lock();
let store = inner.state.downloads();
for id in downloads::withdraw(store, item.0) {
q.pending.push_front((item.0, id));
}
q.priority_active = q.priority_active.saturating_sub(1);
drop(q);
inner.has_work.notify_one();
}
}
}
}
fn run_download(inner: &Arc<Inner>, db_id: i64) {
let store = inner.state.downloads();
let cfg = config::Config::cached();
let result = match crate::helpers::subsonic_client(&cfg) {
None => Some(Err(crate::helpers::remote_unavailable(&cfg))),
Some(client) => {
let checked = std::cell::Cell::new(std::time::Instant::now());
let cancelled = || {
if checked.get().elapsed() < CANCEL_CHECK {
return false;
}
checked.set(std::time::Instant::now());
store.abandoned(db_id)
};
std::panic::catch_unwind(AssertUnwindSafe(|| {
download_track(
db_id,
&cancelled,
&inner.cmd_tx,
&inner.state,
&cfg,
&client,
)
}))
.unwrap_or_else(|_| {
log::error!("download panicked for track {db_id}");
Some(Err("download panicked".into()))
})
}
};
if let Some(Ok(_)) = &result
&& store.kept(db_id)
{
pin(db_id);
}
let mut q = inner.queue.lock();
let settled = match result {
Some(result) => Some(downloads::settle(&inner.state, db_id, &result)),
None => {
for id in downloads::withdraw(store, db_id) {
q.pending.push_front((db_id, id));
}
None
}
};
drop(q);
if let Some(settled) = settled {
settled.announce(&inner.cmd_tx);
}
inner.has_work.notify_all();
if config::Config::cached().cache_limit_bytes().is_some() {
sync(inner);
}
trim_cache(inner, false);
}
fn pin(db_id: i64) {
let pinned = crate::db::pool::shared()
.get()
.map_err(|e| e.to_string())
.and_then(|db| queries::pin_cached(&db.conn, &[db_id]).map_err(|e| e.to_string()));
if let Err(e) = pinned {
log::warn!("could not pin the download of track {db_id}: {e}");
}
}
const EVICT_EVERY: std::time::Duration = std::time::Duration::from_secs(60);
fn trim_cache(inner: &Inner, now: bool) {
{
let mut last = inner.last_evicted.lock();
if !now && last.is_some_and(|t| t.elapsed() < EVICT_EVERY) {
return;
}
*last = Some(std::time::Instant::now());
}
let cfg = config::Config::cached();
if cfg.cache_limit_bytes().is_none() {
return;
}
let keep = inner.queue.lock().window.clone().unwrap_or_else(|| {
inner
.state
.playback_order()
.into_iter()
.map(|(track, _)| track)
.collect()
});
match crate::db::pool::shared().get() {
Ok(db) => {
if crate::helpers::evict_cache(&db, &cfg, &keep, false) > 0 {
inner.state.reset_items_with_missing_files();
}
}
Err(e) => log::warn!("cache eviction: could not open the database: {e}"),
}
}
static LIMIT_CHANGES: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
pub fn cache_limit_changed() {
LIMIT_CHANGES.fetch_add(1, std::sync::atomic::Ordering::AcqRel);
crate::signal::engine_changed().bump();
}
fn worker_loop(inner: Arc<Inner>, index: usize) {
loop {
let item = loop {
if let Some(client) = crate::helpers::subsonic_client(&config::Config::cached()) {
client.outage().hold();
}
let cursor = cursor_download(&inner.state);
let mut q = inner.queue.lock();
if index >= workers_allowed() {
inner.has_work.wait(&mut q);
continue;
}
let store = inner.state.downloads();
match next_item(&mut q, store, cursor) {
Some((db_id, entry)) => {
if store.claim(db_id, entry) {
break db_id;
}
}
None => inner.has_work.wait(&mut q),
}
};
run_download(&inner, item);
}
}
fn follow_playlist(inner: Arc<Inner>) {
let changed = crate::signal::engine_changed();
let mut seen = changed.generation();
let mut last_version: Option<u64> = None;
let mut last_cursor: Option<QueueItemId> = None;
let mut last_limit = LIMIT_CHANGES.load(std::sync::atomic::Ordering::Acquire);
loop {
let version = inner.state.pending_version();
let current = inner.state.cursor();
let limit = LIMIT_CHANGES.load(std::sync::atomic::Ordering::Acquire);
if last_version != Some(version) || current != last_cursor || limit != last_limit {
last_version = Some(version);
sync(&inner);
}
if limit != last_limit {
last_limit = limit;
trim_cache(&inner, true);
}
if current != last_cursor {
last_cursor = current;
inner.has_work.notify_all();
if let Some(cursor_id) = current {
promote_cursor(&inner, cursor_id);
}
}
seen = changed.wait(seen);
}
}
fn promote_cursor(inner: &Arc<Inner>, cursor_id: QueueItemId) {
if inner.state.item_state(cursor_id) != Some(ItemState::Pending) {
return;
}
let mut priority_items = Vec::new();
{
let mut q = inner.queue.lock();
if let Some(pos) = q.pending.iter().position(|(_, qid)| *qid == cursor_id) {
priority_items.push(q.pending.remove(pos).expect("position just found"));
if let Some(next) = q.pending.pop_front() {
priority_items.push(next);
}
}
}
for item in priority_items {
dispatch_priority(inner, item);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::player::state::PlaylistItem;
fn qid() -> QueueItemId {
QueueItemId::new()
}
fn waiters(store: &DownloadStore, db_id: i64) -> HashSet<QueueItemId> {
store.waiters(db_id).into_iter().collect()
}
#[test]
fn priority_lane_never_exceeds_its_permits() {
let (mut q, store) = (Queue::default(), DownloadStore::new());
let mut spawned = 0;
for i in 0..500 {
if claim_priority(&mut q, &store, (i, qid())) == Dispatch::Spawn {
spawned += 1;
}
assert!(
q.priority_active <= PRIORITY_PERMITS,
"priority lane over its permit count at iteration {}",
i
);
}
assert_eq!(spawned, PRIORITY_PERMITS, "only permitted claims may spawn");
assert_eq!(
q.pending.len(),
500 - PRIORITY_PERMITS,
"everything else must be queued, not dropped"
);
}
#[test]
fn released_permits_are_reusable() {
let (mut q, store) = (Queue::default(), DownloadStore::new());
assert_eq!(claim_priority(&mut q, &store, (1, qid())), Dispatch::Spawn);
assert_eq!(claim_priority(&mut q, &store, (2, qid())), Dispatch::Spawn);
assert_eq!(
claim_priority(&mut q, &store, (3, qid())),
Dispatch::Requeued
);
let _ = downloads::withdraw(&store, 1);
q.priority_active -= 1;
assert_eq!(claim_priority(&mut q, &store, (4, qid())), Dispatch::Spawn);
assert!(q.priority_active <= PRIORITY_PERMITS);
}
#[test]
fn an_in_flight_track_is_never_claimed_twice() {
let (mut q, store) = (Queue::default(), DownloadStore::new());
let id = qid();
assert_eq!(claim_priority(&mut q, &store, (1, id)), Dispatch::Spawn);
assert_eq!(
claim_priority(&mut q, &store, (1, id)),
Dispatch::AlreadyRunning
);
assert_eq!(q.priority_active, 1);
assert!(
q.pending.is_empty(),
"a duplicate request must not re-queue the track"
);
}
#[test]
fn playing_a_track_again_joins_the_transfer_already_running() {
let (mut q, store) = (Queue::default(), DownloadStore::new());
let (first, again) = (qid(), qid());
assert_eq!(claim_priority(&mut q, &store, (7, first)), Dispatch::Spawn);
assert_eq!(
claim_priority(&mut q, &store, (7, again)),
Dispatch::AlreadyRunning
);
assert_eq!(q.priority_active, 1, "one transfer, not two");
assert!(q.pending.is_empty());
assert_eq!(
waiters(&store, 7),
HashSet::from([first, again]),
"both entries wait on the one transfer"
);
}
#[test]
fn different_tracks_still_run_side_by_side() {
let (mut q, store) = (Queue::default(), DownloadStore::new());
assert_eq!(claim_priority(&mut q, &store, (1, qid())), Dispatch::Spawn);
assert_eq!(claim_priority(&mut q, &store, (2, qid())), Dispatch::Spawn);
assert_eq!(q.priority_active, 2);
}
#[test]
fn requeued_priority_item_goes_to_the_head_of_the_queue() {
let (mut q, store) = (Queue::default(), DownloadStore::new());
q.pending.push_back((9, qid()));
for i in 0..PRIORITY_PERMITS {
claim_priority(&mut q, &store, (i as i64, qid()));
}
let wanted = qid();
assert_eq!(
claim_priority(&mut q, &store, (7, wanted)),
Dispatch::Requeued
);
assert_eq!(q.pending.front().map(|(_, id)| *id), Some(wanted));
}
#[test]
fn claiming_removes_a_duplicate_queue_entry() {
let (mut q, store) = (Queue::default(), DownloadStore::new());
let id = qid();
q.pending.push_back((1, id));
q.pending.push_back((2, qid()));
assert_eq!(claim_priority(&mut q, &store, (1, id)), Dispatch::Spawn);
assert_eq!(
q.pending.len(),
1,
"the pool must not also pick up the claimed track"
);
}
#[test]
fn the_track_under_the_cursor_goes_first_and_goes_alone() {
let (a, b, c) = (qid(), qid(), qid());
let (mut q, store) = (Queue::default(), DownloadStore::new());
q.pending.extend([(1, a), (2, b), (3, c)]);
assert_eq!(next_item(&mut q, &store, Some((3, c))), Some((3, Some(c))));
store.claim(3, Some(c));
assert_eq!(next_item(&mut q, &store, Some((3, c))), None);
assert_eq!(q.pending.len(), 2, "the rest wait their turn");
let _ = downloads::withdraw(&store, 3);
assert_eq!(next_item(&mut q, &store, None), Some((1, Some(a))));
assert_eq!(next_item(&mut q, &store, None), Some((2, Some(b))));
}
#[test]
fn a_cursor_with_nothing_queued_for_it_holds_nothing_up() {
let (a, elsewhere) = (qid(), qid());
let (mut q, store) = (Queue::default(), DownloadStore::new());
q.pending.push_back((1, a));
assert_eq!(
next_item(&mut q, &store, Some((9, elsewhere))),
Some((1, Some(a)))
);
}
#[test]
fn a_track_wanted_only_in_the_cache_waits_behind_the_playlist() {
let a = qid();
let (mut q, store) = (Queue::default(), DownloadStore::new());
q.cache.push_back(5);
q.pending.push_back((1, a));
assert_eq!(next_item(&mut q, &store, None), Some((1, Some(a))));
assert_eq!(next_item(&mut q, &store, None), Some((5, None)));
}
#[test]
fn the_queue_lets_go_of_what_the_playlist_no_longer_holds() {
let (old, kept, waiter) = (qid(), qid(), qid());
let (mut q, store) = (Queue::default(), DownloadStore::new());
q.pending.extend([(1, old), (2, kept)]);
store.claim(3, Some(waiter));
sync_with(&mut q, &store, &[(2, kept)], None);
assert_eq!(q.pending, VecDeque::from([(2, kept)]));
assert!(waiters(&store, 3).is_empty());
assert!(
store.abandoned(3),
"a transfer nothing waits on any more is on its way to being stopped"
);
}
#[test]
fn the_queue_is_in_the_order_it_is_given() {
let (a, b, c) = (qid(), qid(), qid());
let (mut q, store) = (Queue::default(), DownloadStore::new());
q.pending.extend([(1, a), (2, b), (3, c)]);
assert!(
!sync_with(&mut q, &store, &[(3, c), (1, a), (2, b)], None),
"nothing new"
);
assert_eq!(q.pending, VecDeque::from([(3, c), (1, a), (2, b)]));
let d = qid();
assert!(sync_with(
&mut q,
&store,
&[(3, c), (4, d), (1, a), (2, b)],
None
));
assert_eq!(q.pending, VecDeque::from([(3, c), (4, d), (1, a), (2, b)]));
}
#[test]
fn playing_from_the_middle_fetches_from_there_first() {
let (inner, ids, _) = queue_over(&[1, 2, 3, 4, 5]);
inner.state.set_cursor(Some(ids[2].1));
let wanted = inner.state.pending_downloads();
sync_with(
&mut inner.queue.lock(),
inner.state.downloads(),
&wanted,
None,
);
let order: Vec<i64> = inner.queue.lock().pending.iter().map(|(t, _)| *t).collect();
assert_eq!(order, vec![3, 4, 5, 1, 2]);
}
#[test]
fn a_new_entry_for_a_track_in_flight_waits_on_that_transfer() {
let (running, again) = (qid(), qid());
let (mut q, store) = (Queue::default(), DownloadStore::new());
store.claim(7, Some(running));
assert!(!sync_with(
&mut q,
&store,
&[(7, running), (7, again)],
None
));
assert!(q.pending.is_empty());
assert_eq!(waiters(&store, 7), HashSet::from([running, again]));
}
#[test]
fn a_queue_replaced_with_the_same_track_keeps_its_transfer() {
let (old, new) = (qid(), qid());
let (mut q, store) = (Queue::default(), DownloadStore::new());
store.claim(7, Some(old));
sync_with(&mut q, &store, &[(7, new)], None);
assert_eq!(waiters(&store, 7), HashSet::from([new]));
assert!(!store.abandoned(7));
}
#[test]
fn entries_past_the_window_wait_until_it_reaches_them() {
let (played, next, beyond) = (qid(), qid(), qid());
let (mut q, store) = (Queue::default(), DownloadStore::new());
let wanted = [(2, next), (3, beyond), (1, played)];
sync_with(&mut q, &store, &wanted, Some(&HashSet::from([next])));
assert_eq!(q.pending, VecDeque::from([(2, next)]));
sync_with(&mut q, &store, &wanted, Some(&HashSet::from([beyond])));
assert_eq!(q.pending, VecDeque::from([(3, beyond)]));
}
#[test]
fn a_track_inside_the_window_keeps_its_transfer_when_the_window_is_recomputed() {
let (playing, next, beyond) = (qid(), qid(), qid());
let (mut q, store) = (Queue::default(), DownloadStore::new());
store.claim(2, Some(next));
store.claim(3, Some(beyond));
let wanted = [(1, playing), (2, next), (3, beyond)];
sync_with(
&mut q,
&store,
&wanted,
Some(&HashSet::from([playing, next])),
);
assert!(!store.abandoned(2), "inside the window: still wanted");
assert_eq!(waiters(&store, 2), HashSet::from([next]));
assert!(store.abandoned(3), "past the window: let go");
assert_eq!(q.pending, VecDeque::from([(1, playing)]));
}
fn queue_over(
tracks: &[i64],
) -> (
Arc<Inner>,
Vec<(i64, QueueItemId)>,
crossbeam_channel::Receiver<PlayerCommand>,
) {
crate::config::isolate_config_for_tests();
let item = |db_id: i64| PlaylistItem {
playlist_entry_id: None,
id: qid(),
db_id: Some(db_id),
path: std::path::PathBuf::from(format!("/cache/track-{db_id}.flac")),
title: format!("track-{db_id}"),
artist: "Artist".into(),
album_artist: "Artist".into(),
album: "Album".into(),
year: None,
codec: None,
track_number: None,
disc: None,
duration_ms: None,
state: ItemState::Pending,
played: false,
};
let items: Vec<_> = tracks.iter().map(|&db_id| item(db_id)).collect();
let ids: Vec<_> = items.iter().map(|i| (i.db_id.unwrap(), i.id)).collect();
let state = SharedPlayerState::new();
state.add_items(items);
state.set_cursor(Some(ids[0].1));
let (cmd_tx, cmd_rx) = crossbeam_channel::unbounded();
let inner = Arc::new(Inner {
queue: Mutex::new(Queue::default()),
has_work: Condvar::new(),
state,
cmd_tx,
last_evicted: Mutex::new(None),
spawned: std::sync::atomic::AtomicUsize::new(0),
});
(inner, ids, cmd_rx)
}
fn first_play() -> (Arc<Inner>, Vec<(i64, QueueItemId)>) {
let (inner, ids, _) = queue_over(&[1, 2, 3]);
inner.queue.lock().pending.extend(ids.iter().copied());
(inner, ids)
}
#[test]
fn a_cursor_set_before_its_tracks_were_queued_still_goes_first() {
let (inner, ids) = first_play();
promote_cursor(&inner, ids[0].1);
let left: Vec<_> = inner.queue.lock().pending.iter().copied().collect();
assert_eq!(
left,
vec![ids[2]],
"the cursor's track and the next went to the priority lane"
);
}
fn failed(inner: &Inner, id: QueueItemId) -> bool {
matches!(inner.state.item_state(id), Some(ItemState::Failed(_)))
}
#[test]
fn a_watcher_started_over_a_playlist_fetches_it_unprompted() {
let (inner, ids, _) = queue_over(&[1, 2, 3]);
let watched = inner.clone();
std::thread::spawn(move || follow_playlist(watched));
let deadline = std::time::Instant::now() + Duration::from_secs(5);
while !ids.iter().all(|(_, id)| failed(&inner, *id)) && std::time::Instant::now() < deadline
{
std::thread::sleep(Duration::from_millis(10));
}
assert!(ids.iter().all(|(_, id)| failed(&inner, *id)));
}
fn duplicate_that_fails() -> (
Arc<Inner>,
crossbeam_channel::Receiver<PlayerCommand>,
QueueItemId,
QueueItemId,
) {
let (inner, ids, cmd_rx) = queue_over(&[1, 1]);
let (first, again) = (ids[0].1, ids[1].1);
inner.state.set_cursor(Some(again));
let store = inner.state.downloads();
store.claim(1, Some(first));
store.join(1, Some(again));
(inner, cmd_rx, first, again)
}
#[test]
fn a_failed_transfer_fails_every_entry_waiting_on_it() {
let (inner, cmd_rx, first, again) = duplicate_that_fails();
run_download(&inner, 1);
for id in [first, again] {
assert!(failed(&inner, id), "every waiter hears the one answer");
}
let sent: Vec<_> = cmd_rx.try_iter().collect();
assert!(
matches!(sent.as_slice(), [PlayerCommand::TrackFailed(id)] if *id == again),
"the cursor's entry is told it failed, not that it is ready: {sent:?}"
);
assert!(!inner.state.downloads().in_flight(1));
}
#[test]
fn a_transfer_whose_first_entry_was_removed_still_answers_the_rest() {
let (inner, _cmd_rx, first, again) = duplicate_that_fails();
inner.state.remove_item(first);
sync(&inner);
run_download(&inner, 1);
assert!(failed(&inner, again));
}
#[test]
fn a_transfer_withdrawn_with_an_entry_waiting_queues_it_again() {
let (inner, ids, _) = queue_over(&[1]);
let store = inner.state.downloads();
store.claim(1, Some(ids[0].1));
let mut q = inner.queue.lock();
for id in downloads::withdraw(store, 1) {
q.pending.push_front((1, id));
}
assert_eq!(q.pending, VecDeque::from([ids[0]]));
}
}