use std::collections::{HashMap, HashSet, VecDeque};
use std::panic::AssertUnwindSafe;
use std::sync::{Arc, Mutex as StdMutex};
use parking_lot::{Condvar, Mutex};
use crate::config;
use crate::player::commands::PlayerCommand;
use crate::player::state::{LoadState, QueueItemId, SharedPlayerState};
use crate::helpers::download_track;
const PRIORITY_PERMITS: usize = 2;
#[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>,
log_buf: Arc<StdMutex<Vec<String>>>,
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 || 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, cursor: Option<(i64, QueueItemId)>) -> Option<(i64, QueueItemId)> {
match cursor {
Some((db_id, _)) if q.in_flight.contains_key(&db_id) => 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(),
}
}
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 !matches!(item.state, crate::player::state::ItemState::Pending) {
return None;
}
Some((item.db_id?, id))
}
#[derive(Default)]
struct Queue {
pending: VecDeque<(i64, QueueItemId)>,
in_flight: HashMap<i64, HashSet<QueueItemId>>,
priority_active: usize,
}
#[derive(Debug, PartialEq, Eq)]
enum Dispatch {
Spawn,
Requeued,
AlreadyRunning,
}
fn claim_priority(q: &mut Queue, item: (i64, QueueItemId)) -> Dispatch {
let (db_id, queue_id) = item;
q.pending.retain(|(_, qid)| *qid != queue_id);
if let Some(waiting) = q.in_flight.get_mut(&db_id) {
waiting.insert(queue_id);
return Dispatch::AlreadyRunning;
}
if q.priority_active >= PRIORITY_PERMITS {
q.pending.push_front(item);
return Dispatch::Requeued;
}
q.priority_active += 1;
q.in_flight.insert(db_id, HashSet::from([queue_id]));
Dispatch::Spawn
}
fn release_priority(q: &mut Queue, db_id: i64) {
q.in_flight.remove(&db_id);
q.priority_active = q.priority_active.saturating_sub(1);
}
struct Claim {
inner: Arc<Inner>,
db_id: i64,
priority: bool,
}
impl Drop for Claim {
fn drop(&mut self) {
let mut q = self.inner.queue.lock();
if self.priority {
release_priority(&mut q, self.db_id);
} else {
q.in_flight.remove(&self.db_id);
}
drop(q);
self.inner.has_work.notify_all();
}
}
impl DownloadQueue {
pub fn spawn(
cmd_tx: crossbeam_channel::Sender<PlayerCommand>,
state: Arc<SharedPlayerState>,
log_buf: Arc<StdMutex<Vec<String>>>,
) -> Self {
let inner = Arc::new(Inner {
queue: Mutex::new(Queue::default()),
has_work: Condvar::new(),
state,
cmd_tx,
log_buf,
last_evicted: Mutex::new(None),
spawned: std::sync::atomic::AtomicUsize::new(0),
});
ensure_workers(&inner);
let trimmer = inner.clone();
let _ = std::thread::Builder::new()
.name("koan-dl-trim".into())
.spawn(move || trim_cache(&trimmer));
let watcher_inner = inner.clone();
if let Err(e) = std::thread::Builder::new()
.name("koan-dl-watch".into())
.spawn(move || cursor_watcher(watcher_inner))
{
log::error!("failed to spawn download cursor watcher: {}", e);
}
Self { inner }
}
pub fn enqueue(&self, items: Vec<(i64, QueueItemId)>) {
if items.is_empty() {
return;
}
ensure_workers(&self.inner);
retry_server_now();
self.inner.queue.lock().pending.extend(items);
self.inner.has_work.notify_all();
if let Some(cursor) = self.inner.state.cursor() {
promote_cursor(&self.inner, cursor);
}
}
pub fn prioritize(&self, db_id: i64, queue_id: QueueItemId) {
retry_server_now();
dispatch_priority(&self.inner, (db_id, queue_id));
let album_mates = self.inner.state.same_album_item_ids(queue_id);
if !album_mates.is_empty() {
let mate_set: HashSet<QueueItemId> = album_mates.into_iter().collect();
bump_to_front(&mut self.inner.queue.lock().pending, &mate_set);
self.inner.has_work.notify_all();
}
}
}
fn bump_to_front(pending: &mut VecDeque<(i64, QueueItemId)>, ids: &HashSet<QueueItemId>) {
let (front, rest): (VecDeque<_>, VecDeque<_>) =
pending.drain(..).partition(|(_, qid)| ids.contains(qid));
*pending = front;
pending.extend(rest);
}
fn dispatch_priority(inner: &Arc<Inner>, item: (i64, QueueItemId)) {
let dispatch = claim_priority(&mut inner.queue.lock(), 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 _claim = Claim {
inner: spawn_inner.clone(),
db_id: item.0,
priority: true,
};
run_download(&spawn_inner, item);
});
if let Err(e) = spawned {
log::error!("failed to spawn priority download: {}", e);
let mut q = inner.queue.lock();
release_priority(&mut q, item.0);
q.pending.push_front(item);
drop(q);
inner.has_work.notify_one();
}
}
}
}
fn settle_waiters(inner: &Arc<Inner>, db_id: i64, downloaded: QueueItemId) {
let waiting: Vec<QueueItemId> = {
let q = inner.queue.lock();
q.in_flight
.get(&db_id)
.map(|ids| ids.iter().copied().filter(|id| *id != downloaded).collect())
.unwrap_or_default()
};
if waiting.is_empty() {
return;
}
let Some(item) = inner.state.get_item(downloaded) else {
return;
};
for id in waiting {
inner.state.update_paths(&[(id, item.path.clone())]);
inner.state.update_item_state(id, item.state.clone());
if inner.state.is_cursor(id) {
inner.cmd_tx.send(PlayerCommand::TrackReady(id)).ok();
}
}
}
fn run_download(inner: &Arc<Inner>, (db_id, queue_id): (i64, QueueItemId)) {
let cfg = config::Config::cached();
let Some(client) = crate::helpers::subsonic_client(&cfg) else {
crate::helpers::fail_track(
&inner.state,
&inner.cmd_tx,
queue_id,
crate::helpers::remote_unavailable(&cfg),
);
return;
};
let outcome = std::panic::catch_unwind(AssertUnwindSafe(|| {
download_track(
db_id,
queue_id,
&inner.cmd_tx,
&inner.log_buf,
&inner.state,
&cfg,
&client,
);
}));
if outcome.is_err() {
log::error!("download panicked for {:?}", queue_id);
crate::helpers::fail_track(
&inner.state,
&inner.cmd_tx,
queue_id,
"download panicked".into(),
);
}
settle_waiters(inner, db_id, queue_id);
trim_cache(inner);
}
const EVICT_EVERY: std::time::Duration = std::time::Duration::from_secs(60);
fn trim_cache(inner: &Inner) {
{
let mut last = inner.last_evicted.lock();
if 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
.state
.snapshot_playlist()
.0
.iter()
.filter_map(|i| i.db_id)
.collect();
match crate::db::pool::shared().get() {
Ok(db) => {
crate::helpers::evict_cache(&db, &cfg, &keep, false);
}
Err(e) => log::warn!("cache eviction: could not open the database: {e}"),
}
}
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;
}
match next_item(&mut q, cursor) {
Some(item) => {
match q.in_flight.get_mut(&item.0) {
Some(waiting) => {
waiting.insert(item.1);
}
None => {
q.in_flight.insert(item.0, HashSet::from([item.1]));
break item;
}
}
}
None => inner.has_work.wait(&mut q),
}
};
let _claim = Claim {
inner: inner.clone(),
db_id: item.0,
priority: false,
};
run_download(&inner, item);
}
}
fn cursor_watcher(inner: Arc<Inner>) {
let changed = crate::signal::engine_changed();
let mut seen = changed.generation();
let mut last_cursor: Option<QueueItemId> = None;
loop {
seen = changed.wait(seen);
let current = inner.state.cursor();
if current == last_cursor {
continue;
}
last_cursor = current;
inner.has_work.notify_all();
if let Some(cursor_id) = current {
promote_cursor(&inner, cursor_id);
}
}
}
fn promote_cursor(inner: &Arc<Inner>, cursor_id: QueueItemId) {
let is_pending = inner
.state
.item_load_state(cursor_id)
.is_some_and(|s| matches!(s, LoadState::Pending));
if !is_pending {
return;
}
let album_mate_ids: HashSet<QueueItemId> = inner
.state
.same_album_item_ids(cursor_id)
.into_iter()
.collect();
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 !album_mate_ids.is_empty() {
bump_to_front(&mut q.pending, &album_mate_ids);
}
if let Some(next) = q.pending.pop_front() {
priority_items.push(next);
}
}
}
for item in priority_items {
dispatch_priority(inner, item);
}
}
pub fn shared(
cmd_tx: &crossbeam_channel::Sender<PlayerCommand>,
state: &Arc<SharedPlayerState>,
log_buf: Option<Arc<StdMutex<Vec<String>>>>,
) -> &'static DownloadQueue {
static QUEUE: std::sync::OnceLock<DownloadQueue> = std::sync::OnceLock::new();
QUEUE.get_or_init(|| {
DownloadQueue::spawn(
cmd_tx.clone(),
state.clone(),
log_buf.unwrap_or_else(|| Arc::new(StdMutex::new(Vec::new()))),
)
})
}
#[cfg(test)]
mod tests {
use super::*;
fn qid() -> QueueItemId {
QueueItemId::new()
}
#[test]
fn priority_lane_never_exceeds_its_permits() {
let mut q = Queue::default();
let mut spawned = 0;
for i in 0..500 {
if claim_priority(&mut q, (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 = Queue::default();
assert_eq!(claim_priority(&mut q, (1, qid())), Dispatch::Spawn);
assert_eq!(claim_priority(&mut q, (2, qid())), Dispatch::Spawn);
assert_eq!(claim_priority(&mut q, (3, qid())), Dispatch::Requeued);
release_priority(&mut q, 1);
assert_eq!(claim_priority(&mut q, (4, qid())), Dispatch::Spawn);
assert!(q.priority_active <= PRIORITY_PERMITS);
}
#[test]
fn an_in_flight_track_is_never_claimed_twice() {
let mut q = Queue::default();
let id = qid();
assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::Spawn);
assert_eq!(claim_priority(&mut q, (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 = Queue::default();
let (first, again) = (qid(), qid());
assert_eq!(claim_priority(&mut q, (7, first)), Dispatch::Spawn);
assert_eq!(claim_priority(&mut q, (7, again)), Dispatch::AlreadyRunning);
assert_eq!(q.priority_active, 1, "one transfer, not two");
assert!(q.pending.is_empty());
assert_eq!(
q.in_flight.get(&7),
Some(&HashSet::from([first, again])),
"both entries wait on the one transfer"
);
}
#[test]
fn a_worker_picking_up_a_duplicate_waits_on_the_running_one() {
let mut q = Queue::default();
let (running, queued) = (qid(), qid());
assert_eq!(claim_priority(&mut q, (7, running)), Dispatch::Spawn);
match q.in_flight.get_mut(&7) {
Some(waiting) => {
waiting.insert(queued);
}
None => panic!("the track should already be claimed"),
}
assert_eq!(
q.in_flight.get(&7),
Some(&HashSet::from([running, queued])),
"the queued entry waits rather than starting a second transfer"
);
}
#[test]
fn different_tracks_still_run_side_by_side() {
let mut q = Queue::default();
assert_eq!(claim_priority(&mut q, (1, qid())), Dispatch::Spawn);
assert_eq!(claim_priority(&mut q, (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 = Queue::default();
q.pending.push_back((9, qid()));
for i in 0..PRIORITY_PERMITS {
claim_priority(&mut q, (i as i64, qid()));
}
let wanted = qid();
assert_eq!(claim_priority(&mut q, (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 = Queue::default();
let id = qid();
q.pending.push_back((1, id));
q.pending.push_back((2, qid()));
assert_eq!(claim_priority(&mut q, (1, id)), Dispatch::Spawn);
assert_eq!(
q.pending.len(),
1,
"the pool must not also pick up the claimed track"
);
}
#[test]
fn bump_to_front_preserves_relative_order() {
let (a, b, c, d) = (qid(), qid(), qid(), qid());
let mut pending: VecDeque<(i64, QueueItemId)> =
[(1, a), (2, b), (3, c), (4, d)].into_iter().collect();
let mates: HashSet<QueueItemId> = [b, d].into_iter().collect();
bump_to_front(&mut pending, &mates);
let order: Vec<QueueItemId> = pending.iter().map(|(_, id)| *id).collect();
assert_eq!(order, vec![b, d, a, c]);
}
#[test]
fn the_track_under_the_cursor_goes_first_and_goes_alone() {
let (a, b, c) = (qid(), qid(), qid());
let mut q = Queue::default();
q.pending.extend([(1, a), (2, b), (3, c)]);
assert_eq!(next_item(&mut q, Some((3, c))), Some((3, c)));
q.in_flight.insert(3, HashSet::from([c]));
assert_eq!(next_item(&mut q, Some((3, c))), None);
assert_eq!(q.pending.len(), 2, "the rest wait their turn");
q.in_flight.remove(&3);
assert_eq!(next_item(&mut q, None), Some((1, a)));
assert_eq!(next_item(&mut q, None), Some((2, b)));
}
#[test]
fn a_cursor_with_nothing_queued_for_it_holds_nothing_up() {
let (a, elsewhere) = (qid(), qid());
let mut q = Queue::default();
q.pending.push_back((1, a));
assert_eq!(next_item(&mut q, Some((9, elsewhere))), Some((1, a)));
}
#[test]
fn a_cursor_set_before_its_tracks_were_queued_still_goes_first() {
crate::config::isolate_config_for_tests();
let item = |title: &str, db_id: i64| crate::player::state::PlaylistItem {
playlist_entry_id: None,
id: qid(),
db_id: Some(db_id),
path: std::path::PathBuf::from(format!("/cache/{title}.flac")),
title: title.into(),
artist: "Artist".into(),
album_artist: "Artist".into(),
album: "Album".into(),
year: None,
codec: None,
track_number: None,
disc: None,
duration_ms: None,
state: crate::player::state::ItemState::Pending,
};
let items = vec![item("one", 1), item("two", 2), item("three", 3)];
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,
log_buf: Arc::new(StdMutex::new(Vec::new())),
last_evicted: Mutex::new(None),
spawned: std::sync::atomic::AtomicUsize::new(0),
});
inner.queue.lock().pending.extend(ids.iter().copied());
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"
);
}
}