use std::{
sync::{
atomic::Ordering,
mpsc::{self, Receiver, Sender},
},
thread,
};
use indicatrix::optics::chromophore::{ChromophoreCatalogue, SolveResult};
use indicatrix_cut_core::material::color::solve_cancellable;
use super::physics_state::SolveJob;
#[derive(Debug, Clone)]
pub struct SolveOutcome {
pub generation: u64,
pub result: SolveResult,
}
pub struct SolverWorker {
tx: Sender<SolveJob>,
}
fn solve_real(job: &SolveJob) -> Option<SolveResult> {
solve_cancellable(
ChromophoreCatalogue::global(),
&job.host,
job.target_lab,
job.reference_path_mm,
&job.locked,
&job.cancel,
)
}
impl SolverWorker {
#[must_use]
pub fn spawn(on_done: impl Fn(SolveOutcome) + Send + 'static) -> Self {
Self::spawn_with(solve_real, on_done)
}
#[must_use]
pub fn spawn_with(
solve: impl Fn(&SolveJob) -> Option<SolveResult> + Send + 'static,
on_done: impl Fn(SolveOutcome) + Send + 'static,
) -> Self {
let (tx, rx) = mpsc::channel::<SolveJob>();
let _ = thread::Builder::new()
.name("physics-color-solver".to_string())
.spawn(move || run(&rx, &solve, &on_done));
Self { tx }
}
pub fn submit(&self, job: SolveJob) {
let _ = self.tx.send(job);
}
}
fn run(
rx: &Receiver<SolveJob>,
solve: &impl Fn(&SolveJob) -> Option<SolveResult>,
on_done: &impl Fn(SolveOutcome),
) {
while let Ok(mut job) = rx.recv() {
while let Ok(newer) = rx.try_recv() {
job = newer;
}
if job.cancel.load(Ordering::SeqCst) {
continue;
}
let Some(result) = solve(&job) else {
continue;
};
if job.cancel.load(Ordering::SeqCst) {
continue;
}
on_done(SolveOutcome {
generation: job.generation,
result,
});
}
}
#[cfg(test)]
mod tests {
use super::*;
use indicatrix::{color::body_color::BodyColor, optics::chromophore::ColorRecipe};
use std::{
sync::{Arc, Mutex, atomic::AtomicBool},
time::Duration,
};
fn job(generation: u64) -> SolveJob {
SolveJob {
generation,
host: "h".to_string(),
target_lab: [50.0, 0.0, 0.0],
reference_path_mm: 5.0,
locked: Vec::new(),
cancel: Arc::new(AtomicBool::new(false)),
}
}
fn fake_result(tag: &str) -> SolveResult {
SolveResult {
recipe: ColorRecipe::new(tag, 1),
achieved: BodyColor {
xyz: [0.0; 3],
lab: [0.0; 3],
srgb: [0; 3],
out_of_gamut: false,
},
delta_e: 0.0,
delta_e_a: None,
reachable: true,
capped: false,
evals: 0,
}
}
#[test]
fn the_latest_request_wins_and_stale_results_are_dropped() {
let (gate_tx, gate_rx) = mpsc::channel::<()>();
let gate_rx = Mutex::new(gate_rx);
let (done_tx, done_rx) = mpsc::channel::<u64>();
let done_tx = Mutex::new(done_tx);
let worker = SolverWorker::spawn_with(
move |j| {
if j.generation == 1 {
let _ = gate_rx.lock().unwrap().recv_timeout(Duration::from_secs(5));
}
Some(fake_result(&j.host))
},
move |outcome| {
let _ = done_tx.lock().unwrap().send(outcome.generation);
},
);
let first = job(1);
let first_cancel = Arc::clone(&first.cancel);
let started = std::time::Instant::now();
worker.submit(first);
std::thread::sleep(Duration::from_millis(50));
first_cancel.store(true, Ordering::SeqCst);
worker.submit(job(2));
worker.submit(job(3));
assert!(
started.elapsed() < Duration::from_secs(1),
"submit must not block"
);
gate_tx.send(()).unwrap();
let mut got = Vec::new();
while let Ok(g) = done_rx.recv_timeout(Duration::from_millis(500)) {
got.push(g);
}
assert_eq!(got.last(), Some(&3), "the latest request wins: {got:?}");
assert!(!got.contains(&1), "the cancelled solve is dropped: {got:?}");
assert!(
!got.contains(&2),
"the coalesced request never runs: {got:?}"
);
}
}