use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::mpsc::{Receiver, Sender, TryRecvError, channel};
use std::time::{Duration, Instant};
use crate::acquire::Registry;
use crate::acquire::shop::{self, GroupOutcome, QuerySpec, SearchOpts};
use crate::acquire::types::{AcquiredFile, AudioFormat, FetchOpts, ItemRef, Retention};
use crate::config::{Config, Credentials};
const FETCH_BUDGET: Duration = Duration::from_secs(1800);
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
pub enum JobKind {
Search,
Fetch,
Probe,
Fingerprint,
}
impl JobKind {
pub fn label(self) -> &'static str {
match self {
Self::Search => "search",
Self::Fetch => "download",
Self::Probe => "file check",
Self::Fingerprint => "fingerprint",
}
}
}
pub enum Job {
Shop {
specs: Vec<QuerySpec>,
opts: Box<SearchOpts>,
},
Fetch {
item: ItemRef,
dest: PathBuf,
format_pref: Vec<AudioFormat>,
overwrite: bool,
},
Probe {
entry_id: i64,
generation: u64,
path: PathBuf,
},
Fingerprint {
entry_id: i64,
generation: u64,
src: Box<crate::analysis::TrackHeader>,
dst_path: PathBuf,
dst_length: Option<i64>,
dst_bpm: Option<i64>,
},
}
impl Job {
pub fn kind(&self) -> JobKind {
match self {
Self::Shop { .. } => JobKind::Search,
Self::Fetch { .. } => JobKind::Fetch,
Self::Probe { .. } => JobKind::Probe,
Self::Fingerprint { .. } => JobKind::Fingerprint,
}
}
}
pub enum Update {
Started,
Progress {
done: usize,
total: usize,
label: String,
},
Note(String),
Finished(Box<Vec<GroupOutcome>>),
Fetched(Box<Result<Vec<AcquiredFile>, String>>),
Probed(Box<(i64, u64, Result<crate::audio::AudioInfo, String>)>),
Fingerprinted(Box<(i64, u64, Result<crate::transfer::GateOutcome, String>)>),
Failed(String),
}
impl Update {
fn is_terminal(&self) -> bool {
matches!(
self,
Self::Finished(_)
| Self::Fetched(_)
| Self::Probed(_)
| Self::Fingerprinted(_)
| Self::Failed(_)
)
}
fn ends(&self) -> Option<JobKind> {
match self {
Self::Finished(_) => Some(JobKind::Search),
Self::Fetched(_) => Some(JobKind::Fetch),
Self::Probed(_) => Some(JobKind::Probe),
Self::Fingerprinted(_) => Some(JobKind::Fingerprint),
Self::Failed(_) | Self::Started | Self::Progress { .. } | Self::Note(_) => None,
}
}
}
pub struct Worker {
jobs: Option<Sender<Job>>,
updates: Receiver<Update>,
outstanding: usize,
by_kind: HashMap<JobKind, usize>,
}
impl Worker {
pub fn spawn(cfg: &Config, creds: &Credentials) -> std::io::Result<Self> {
let (job_tx, job_rx) = channel::<Job>();
let (up_tx, up_rx) = channel::<Update>();
let cfg = cfg.clone();
let creds = creds.clone();
std::thread::Builder::new()
.name("rr-shop".into())
.spawn(move || {
let notes = up_tx.clone();
crate::acquire::route_progress(Box::new(move |line| {
let _ = notes.send(Update::Note(line.to_string()));
}));
let reg = Registry::from_config(&cfg, &creds);
while let Ok(job) = job_rx.recv() {
match job {
Job::Shop { specs, opts } => {
if up_tx.send(Update::Started).is_err() {
return;
}
let tx = up_tx.clone();
let groups =
shop::search_many(®, &specs, &opts, |done, total, label| {
let _ = tx.send(Update::Progress {
done,
total,
label: label.to_string(),
});
});
if up_tx.send(Update::Finished(Box::new(groups))).is_err() {
return;
}
}
Job::Fetch {
item,
dest,
format_pref,
overwrite,
} => {
if up_tx.send(Update::Started).is_err() {
return;
}
let result = match reg.get(item.backend) {
None => Err(format!("{} is not enabled", item.backend)),
Some(backend) => backend
.fetch(
&item,
&FetchOpts {
dest_dir: dest,
format_pref,
retention: Retention::Keep,
overwrite,
deadline: Instant::now() + FETCH_BUDGET,
},
)
.map_err(|e| e.to_string()),
};
if up_tx.send(Update::Fetched(Box::new(result))).is_err() {
return;
}
}
Job::Probe {
entry_id,
generation,
path,
} => {
if up_tx.send(Update::Started).is_err() {
return;
}
let result = crate::audio::probe(&path).map_err(|e| e.to_string());
let msg = Update::Probed(Box::new((entry_id, generation, result)));
if up_tx.send(msg).is_err() {
return;
}
}
Job::Fingerprint {
entry_id,
generation,
src,
dst_path,
dst_length,
dst_bpm,
} => {
if up_tx.send(Update::Started).is_err() {
return;
}
let result =
crate::transfer::gate(&src, &dst_path, dst_length, dst_bpm, &cfg)
.map_err(|e| e.to_string());
let msg =
Update::Fingerprinted(Box::new((entry_id, generation, result)));
if up_tx.send(msg).is_err() {
return;
}
}
}
}
})?;
Ok(Self {
jobs: Some(job_tx),
updates: up_rx,
outstanding: 0,
by_kind: HashMap::new(),
})
}
pub fn is_busy(&self) -> bool {
self.outstanding > 0
}
pub fn outstanding(&self) -> usize {
self.outstanding
}
pub fn outstanding_of(&self, kind: JobKind) -> usize {
self.by_kind.get(&kind).copied().unwrap_or(0)
}
pub fn submit(&mut self, job: Job) -> bool {
let kind = job.kind();
match self.jobs.as_ref().map(|tx| tx.send(job)) {
Some(Ok(())) => {
self.outstanding += 1;
*self.by_kind.entry(kind).or_insert(0) += 1;
true
}
_ => false,
}
}
pub fn drain(&mut self) -> Vec<Update> {
let mut out = Vec::new();
loop {
match self.updates.try_recv() {
Ok(u) => {
if u.is_terminal() {
self.outstanding = self.outstanding.saturating_sub(1);
}
if let Some(kind) = u.ends()
&& let Some(n) = self.by_kind.get_mut(&kind)
{
*n = n.saturating_sub(1);
}
out.push(u);
}
Err(TryRecvError::Empty) => break,
Err(TryRecvError::Disconnected) => {
while self.outstanding > 0 {
self.outstanding -= 1;
out.push(Update::Failed("the worker thread stopped".into()));
}
self.by_kind.clear();
break;
}
}
}
out
}
pub fn shutdown(&mut self) {
self.jobs = None;
}
}
impl Drop for Worker {
fn drop(&mut self) {
self.shutdown();
}
}
#[cfg(test)]
mod tests {
use super::*;
fn worker() -> Worker {
Worker::spawn(&Config::default(), &Credentials::default()).unwrap()
}
fn empty_spec() -> QuerySpec {
QuerySpec {
label: "nothing".into(),
src_id: None,
query: crate::acquire::types::SearchQuery::from_text("", 1),
}
}
#[test]
fn a_new_worker_is_idle_and_has_nothing_to_report() {
let mut w = worker();
assert!(!w.is_busy());
assert!(
w.drain().is_empty(),
"draining must not block or invent updates"
);
}
#[test]
fn several_jobs_queue_behind_each_other() {
let mut w = worker();
let job = || Job::Shop {
specs: vec![empty_spec()],
opts: Box::new(SearchOpts::default()),
};
assert!(w.submit(job()));
assert!(w.submit(job()));
assert!(w.submit(job()));
assert!(w.is_busy());
assert_eq!(w.outstanding(), 3);
}
#[test]
fn every_job_kind_finishes() {
for job in [
Job::Probe {
entry_id: 1,
generation: 0,
path: PathBuf::from("/nonexistent/rr-test.flac"),
},
Job::Shop {
specs: vec![empty_spec()],
opts: Box::new(SearchOpts::default()),
},
] {
let kind = job.kind();
let mut w = worker();
assert!(w.submit(job));
assert_eq!(w.outstanding_of(kind), 1, "{kind:?} was not counted");
let start = Instant::now();
while w.outstanding() > 0 && start.elapsed() < Duration::from_secs(30) {
w.drain();
std::thread::sleep(Duration::from_millis(10));
}
assert_eq!(w.outstanding(), 0, "{kind:?} never cleared the counter");
assert_eq!(w.outstanding_of(kind), 0, "{kind:?} left its own count");
}
}
#[test]
fn one_kind_of_work_does_not_mask_another() {
let mut w = worker();
assert!(w.submit(Job::Probe {
entry_id: 1,
generation: 0,
path: PathBuf::from("/nonexistent/rr-test.flac"),
}));
assert_eq!(w.outstanding_of(JobKind::Probe), 1);
assert_eq!(w.outstanding_of(JobKind::Search), 0);
assert_eq!(w.outstanding_of(JobKind::Fetch), 0);
}
#[test]
fn a_finished_job_clears_busy_and_reports_an_outcome() {
let mut w = worker();
assert!(w.submit(Job::Shop {
specs: vec![empty_spec()],
opts: Box::new(SearchOpts::default()),
}));
let mut updates = Vec::new();
for _ in 0..200 {
updates.extend(w.drain());
if updates
.iter()
.any(|u| matches!(u, Update::Finished(_) | Update::Failed(_)))
{
break;
}
std::thread::sleep(std::time::Duration::from_millis(25));
}
assert!(
updates.iter().any(|u| matches!(u, Update::Started)),
"the UI needs a Started to show progress"
);
assert!(
updates
.iter()
.any(|u| matches!(u, Update::Finished(_) | Update::Failed(_))),
"never got a terminal update"
);
assert!(!w.is_busy(), "busy must clear so another search can run");
assert_eq!(w.outstanding(), 0);
}
#[test]
fn shutdown_is_idempotent_and_does_not_block() {
let mut w = worker();
w.shutdown();
w.shutdown();
assert!(!w.submit(Job::Shop {
specs: vec![empty_spec()],
opts: Box::new(SearchOpts::default()),
}));
}
#[test]
fn dropping_the_worker_stops_the_thread() {
let before = std::thread::available_parallelism().is_ok();
{
let _w = worker();
}
assert!(before, "sanity");
}
}