| 1 |
|
| 2 |
|
| 3 |
|
| 4 |
|
| 5 |
|
| 6 |
|
| 7 |
use std::path::PathBuf; |
| 8 |
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; |
| 9 |
use std::sync::{mpsc, Arc, Mutex}; |
| 10 |
use std::thread; |
| 11 |
|
| 12 |
use rayon::prelude::*; |
| 13 |
|
| 14 |
use super::config::AnalysisConfig; |
| 15 |
use super::suggest::TagSuggestion; |
| 16 |
use super::{analyze_sample, AnalysisResult}; |
| 17 |
use tracing::instrument; |
| 18 |
|
| 19 |
|
| 20 |
pub enum WorkerCommand { |
| 21 |
|
| 22 |
AnalyzeBatch { |
| 23 |
|
| 24 |
samples: Vec<(String, String, PathBuf)>, |
| 25 |
|
| 26 |
config: AnalysisConfig, |
| 27 |
}, |
| 28 |
|
| 29 |
Cancel, |
| 30 |
|
| 31 |
Shutdown, |
| 32 |
} |
| 33 |
|
| 34 |
|
| 35 |
pub enum WorkerEvent { |
| 36 |
|
| 37 |
Progress { |
| 38 |
|
| 39 |
completed: usize, |
| 40 |
|
| 41 |
total: usize, |
| 42 |
|
| 43 |
current_name: String, |
| 44 |
}, |
| 45 |
|
| 46 |
SampleDone { |
| 47 |
|
| 48 |
result: Box<AnalysisResult>, |
| 49 |
|
| 50 |
suggestions: Vec<TagSuggestion>, |
| 51 |
}, |
| 52 |
|
| 53 |
SampleError { |
| 54 |
|
| 55 |
hash: String, |
| 56 |
|
| 57 |
error: String, |
| 58 |
}, |
| 59 |
|
| 60 |
BatchComplete, |
| 61 |
} |
| 62 |
|
| 63 |
|
| 64 |
|
| 65 |
|
| 66 |
|
| 67 |
|
| 68 |
pub struct WorkerHandle { |
| 69 |
cmd_tx: mpsc::Sender<WorkerCommand>, |
| 70 |
event_rx: Mutex<mpsc::Receiver<WorkerEvent>>, |
| 71 |
cancel_flag: Arc<AtomicBool>, |
| 72 |
thread: Option<thread::JoinHandle<()>>, |
| 73 |
} |
| 74 |
|
| 75 |
impl WorkerHandle { |
| 76 |
|
| 77 |
pub fn try_recv(&self) -> Option<WorkerEvent> { |
| 78 |
self.event_rx.lock().ok()?.try_recv().ok() |
| 79 |
} |
| 80 |
|
| 81 |
|
| 82 |
pub fn send(&self, cmd: WorkerCommand) { |
| 83 |
if matches!(cmd, WorkerCommand::Cancel) { |
| 84 |
self.cancel_flag.store(true, Ordering::Release); |
| 85 |
} |
| 86 |
let _ = self.cmd_tx.send(cmd); |
| 87 |
} |
| 88 |
} |
| 89 |
|
| 90 |
impl Drop for WorkerHandle { |
| 91 |
fn drop(&mut self) { |
| 92 |
self.cancel_flag.store(true, Ordering::Release); |
| 93 |
let _ = self.cmd_tx.send(WorkerCommand::Shutdown); |
| 94 |
if let Some(handle) = self.thread.take() { |
| 95 |
let _ = handle.join(); |
| 96 |
} |
| 97 |
} |
| 98 |
} |
| 99 |
|
| 100 |
|
| 101 |
|
| 102 |
|
| 103 |
|
| 104 |
#[instrument(skip_all)] |
| 105 |
pub fn spawn_worker() -> std::io::Result<WorkerHandle> { |
| 106 |
let (cmd_tx, cmd_rx) = mpsc::channel::<WorkerCommand>(); |
| 107 |
let (event_tx, event_rx) = mpsc::channel::<WorkerEvent>(); |
| 108 |
let cancel_flag = Arc::new(AtomicBool::new(false)); |
| 109 |
let cancel_clone = Arc::clone(&cancel_flag); |
| 110 |
|
| 111 |
let thread = thread::Builder::new() |
| 112 |
.name("analysis-worker".to_string()) |
| 113 |
.spawn(move || { |
| 114 |
worker_loop(cmd_rx, event_tx, cancel_clone); |
| 115 |
})?; |
| 116 |
|
| 117 |
Ok(WorkerHandle { |
| 118 |
cmd_tx, |
| 119 |
event_rx: Mutex::new(event_rx), |
| 120 |
cancel_flag, |
| 121 |
thread: Some(thread), |
| 122 |
}) |
| 123 |
} |
| 124 |
|
| 125 |
fn worker_loop( |
| 126 |
cmd_rx: mpsc::Receiver<WorkerCommand>, |
| 127 |
event_tx: mpsc::Sender<WorkerEvent>, |
| 128 |
cancel_flag: Arc<AtomicBool>, |
| 129 |
) { |
| 130 |
while let Ok(cmd) = cmd_rx.recv() { |
| 131 |
match cmd { |
| 132 |
WorkerCommand::Shutdown => break, |
| 133 |
WorkerCommand::Cancel => { |
| 134 |
|
| 135 |
continue; |
| 136 |
} |
| 137 |
WorkerCommand::AnalyzeBatch { samples, config } => { |
| 138 |
|
| 139 |
cancel_flag.store(false, Ordering::Release); |
| 140 |
|
| 141 |
let total = samples.len(); |
| 142 |
let completed = Arc::new(AtomicUsize::new(0)); |
| 143 |
|
| 144 |
|
| 145 |
samples.par_iter().for_each(|(hash, _ext, path)| { |
| 146 |
|
| 147 |
if cancel_flag.load(Ordering::Acquire) { |
| 148 |
return; |
| 149 |
} |
| 150 |
|
| 151 |
let name = crate::util::get_filename(path, "unknown"); |
| 152 |
let done = completed.load(Ordering::Acquire); |
| 153 |
|
| 154 |
let _ = event_tx.send(WorkerEvent::Progress { |
| 155 |
completed: done, |
| 156 |
total, |
| 157 |
current_name: name, |
| 158 |
}); |
| 159 |
|
| 160 |
|
| 161 |
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { |
| 162 |
analyze_sample(hash, path, &config) |
| 163 |
})); |
| 164 |
|
| 165 |
match result { |
| 166 |
Ok(Ok(result)) => { |
| 167 |
let suggestions = if config.auto_suggest_tags { |
| 168 |
super::suggest::suggest_tags(&result) |
| 169 |
} else { |
| 170 |
Vec::new() |
| 171 |
}; |
| 172 |
let _ = event_tx.send(WorkerEvent::SampleDone { |
| 173 |
result: Box::new(result), |
| 174 |
suggestions, |
| 175 |
}); |
| 176 |
} |
| 177 |
Ok(Err(e)) => { |
| 178 |
let _ = event_tx.send(WorkerEvent::SampleError { |
| 179 |
hash: hash.clone(), |
| 180 |
error: e.to_string(), |
| 181 |
}); |
| 182 |
} |
| 183 |
Err(_panic) => { |
| 184 |
let _ = event_tx.send(WorkerEvent::SampleError { |
| 185 |
hash: hash.clone(), |
| 186 |
error: "analysis panicked (internal error)".to_string(), |
| 187 |
}); |
| 188 |
} |
| 189 |
} |
| 190 |
|
| 191 |
completed.fetch_add(1, Ordering::Release); |
| 192 |
}); |
| 193 |
|
| 194 |
let _ = event_tx.send(WorkerEvent::BatchComplete); |
| 195 |
} |
| 196 |
} |
| 197 |
} |
| 198 |
} |
| 199 |
|
| 200 |
#[cfg(test)] |
| 201 |
mod tests { |
| 202 |
use super::*; |
| 203 |
|
| 204 |
#[test] |
| 205 |
fn worker_command_cancel_construction() { |
| 206 |
|
| 207 |
let cmd = WorkerCommand::Cancel; |
| 208 |
assert!(matches!(cmd, WorkerCommand::Cancel)); |
| 209 |
} |
| 210 |
|
| 211 |
#[test] |
| 212 |
fn worker_command_shutdown_construction() { |
| 213 |
let cmd = WorkerCommand::Shutdown; |
| 214 |
assert!(matches!(cmd, WorkerCommand::Shutdown)); |
| 215 |
} |
| 216 |
|
| 217 |
#[test] |
| 218 |
fn worker_command_analyze_batch_construction() { |
| 219 |
let config = AnalysisConfig::default(); |
| 220 |
let samples = vec![ |
| 221 |
("hash1".to_string(), "wav".to_string(), PathBuf::from("/a.wav")), |
| 222 |
("hash2".to_string(), "mp3".to_string(), PathBuf::from("/b.mp3")), |
| 223 |
]; |
| 224 |
let cmd = WorkerCommand::AnalyzeBatch { |
| 225 |
samples: samples.clone(), |
| 226 |
config, |
| 227 |
}; |
| 228 |
match cmd { |
| 229 |
WorkerCommand::AnalyzeBatch { samples: s, config: c } => { |
| 230 |
assert_eq!(s.len(), 2); |
| 231 |
assert_eq!(s[0].0, "hash1"); |
| 232 |
assert_eq!(s[1].2, PathBuf::from("/b.mp3")); |
| 233 |
assert!(c.loudness); |
| 234 |
} |
| 235 |
_ => panic!("expected AnalyzeBatch"), |
| 236 |
} |
| 237 |
} |
| 238 |
|
| 239 |
#[test] |
| 240 |
fn spawn_worker_and_shutdown_via_drop() { |
| 241 |
|
| 242 |
let handle = spawn_worker().unwrap(); |
| 243 |
|
| 244 |
assert!(handle.try_recv().is_none()); |
| 245 |
|
| 246 |
drop(handle); |
| 247 |
} |
| 248 |
|
| 249 |
#[test] |
| 250 |
fn worker_event_variants_constructible() { |
| 251 |
let _progress = WorkerEvent::Progress { |
| 252 |
completed: 5, |
| 253 |
total: 10, |
| 254 |
current_name: "kick.wav".to_string(), |
| 255 |
}; |
| 256 |
let _sample_err = WorkerEvent::SampleError { |
| 257 |
hash: "abc".to_string(), |
| 258 |
error: "decode failed".to_string(), |
| 259 |
}; |
| 260 |
let _batch_complete = WorkerEvent::BatchComplete; |
| 261 |
|
| 262 |
} |
| 263 |
|
| 264 |
#[test] |
| 265 |
fn cancel_flag_set_on_cancel() { |
| 266 |
let handle = spawn_worker().unwrap(); |
| 267 |
handle.send(WorkerCommand::Cancel); |
| 268 |
|
| 269 |
std::thread::sleep(std::time::Duration::from_millis(10)); |
| 270 |
assert!(handle.cancel_flag.load(Ordering::Acquire)); |
| 271 |
drop(handle); |
| 272 |
} |
| 273 |
} |
| 274 |
|