Allow usage on wasm32-unknown-unknown target (#194)

* Avoid using time API when we don't need it

This avoids a syscall to the time API when the result is ignored later anyway.

This allows to use the library with default options on wasm32-unknown-unknown, where the unimplemented syscall would panic otherwise.

* eval_send doesn't need to be an Option

We can drop the value manually, thus avoiding unwrap on each access.

* Keep single `use rayon::prelude::*`

If either `rayon` is already imported, then `rayon::prelude::*` should always resolve.

* Extract comparator

* Fully enable non-parallel mode
This commit is contained in:
Ingvar Stepanyan 2019-09-25 18:01:53 +01:00 committed by Josh Holmer
parent 68db304a2a
commit 659717cb1f
2 changed files with 120 additions and 82 deletions

View file

@ -11,15 +11,13 @@ use crate::png::STD_STRATEGY;
use crate::png::STD_WINDOW; use crate::png::STD_WINDOW;
#[cfg(not(feature = "parallel"))] #[cfg(not(feature = "parallel"))]
use crate::rayon; use crate::rayon;
#[cfg(not(feature = "parallel"))]
use crate::rayon::prelude::*;
use crate::Deadline; use crate::Deadline;
#[cfg(feature = "parallel")] #[cfg(feature = "parallel")]
use rayon; use rayon;
#[cfg(feature = "parallel")]
use rayon::prelude::*; use rayon::prelude::*;
use std::sync::atomic::AtomicUsize; use std::sync::atomic::AtomicUsize;
use std::sync::atomic::Ordering::SeqCst; use std::sync::atomic::Ordering::SeqCst;
#[cfg(feature = "parallel")]
use std::sync::mpsc::*; use std::sync::mpsc::*;
use std::sync::Arc; use std::sync::Arc;
use std::thread; use std::thread;
@ -35,103 +33,15 @@ struct Candidate {
nth: usize, nth: usize,
} }
/// Collect image versions and pick one that compresses best #[derive(Default)]
pub(crate) struct Evaluator { struct Comparator {
deadline: Arc<Deadline>, best_result: Option<Candidate>,
nth: AtomicUsize,
best_candidate_size: Arc<AtomicMin>,
/// images are sent to the thread for evaluation
eval_send: Option<SyncSender<Candidate>>,
// the thread helps evaluate images asynchronously
eval_thread: thread::JoinHandle<Option<PngData>>,
} }
impl Evaluator { impl Comparator {
pub fn new(deadline: Arc<Deadline>) -> Self { fn evaluate(&mut self, new: Candidate) {
let (tx, rx) = sync_channel(4);
Self {
deadline,
best_candidate_size: Arc::new(AtomicMin::new(None)),
nth: AtomicUsize::new(0),
eval_send: Some(tx),
eval_thread: thread::spawn(move || Self::evaluate_images(rx)),
}
}
/// Wait for all evaluations to finish and return smallest reduction
/// Or `None` if all reductions were worse than baseline.
pub fn get_result(mut self) -> Option<PngData> {
let _ = self.eval_send.take(); // disconnect the sender, breaking the loop in the thread
self.eval_thread.join().expect("eval thread")
}
/// Set baseline image. It will be used only to measure minimum compression level required
pub fn set_baseline(&self, image: Arc<PngImage>) {
self.try_image_inner(image, 1.0, false)
}
/// Check if the image is smaller than others
/// Bias is a value in 0..=1 range. Compressed size is multiplied by
/// this fraction when comparing to the best, so 0.95 allows 5% larger size.
pub fn try_image(&self, image: Arc<PngImage>, bias: f32) {
self.try_image_inner(image, bias, true)
}
fn try_image_inner(&self, image: Arc<PngImage>, bias: f32, is_reduction: bool) {
let nth = self.nth.fetch_add(1, SeqCst);
// These clones are only cheap refcounts
let deadline = self.deadline.clone();
let best_candidate_size = self.best_candidate_size.clone();
// sends it off asynchronously for compression,
// but results will be collected via the message queue
let eval_send = self.eval_send.clone();
rayon::spawn(move || {
let filters_iter = STD_FILTERS.par_iter().with_max_len(1);
// Updating of best result inside the parallel loop would require locks,
// which are dangerous to do in side Rayon's loop.
// Instead, only update (atomic) best size in real time,
// and the best result later without need for locks.
filters_iter.for_each(|&filter| {
if deadline.passed() {
return;
}
if let Ok(idat_data) = deflate::deflate(
&image.filter_image(filter),
STD_COMPRESSION,
STD_STRATEGY,
STD_WINDOW,
&best_candidate_size,
&deadline,
) {
best_candidate_size.set_min(idat_data.len());
// the rest is shipped to the evavluation/collection thread
eval_send
.as_ref()
.expect("not finished yet")
.send(Candidate {
image: PngData {
idat_data,
raw: Arc::clone(&image),
},
bias,
filter,
is_reduction,
nth,
})
.expect("send");
}
});
});
}
/// Main loop of evaluation thread
fn evaluate_images(from_channel: Receiver<Candidate>) -> Option<PngData> {
let mut best_result: Option<Candidate> = None;
// ends when the last sender is dropped
for new in from_channel.iter() {
// a tie-breaker is required to make evaluation deterministic // a tie-breaker is required to make evaluation deterministic
let is_best = if let Some(ref old) = best_result { let is_best = if let Some(ref old) = self.best_result {
// ordering is important - later file gets to use bias over earlier, but not the other way // ordering is important - later file gets to use bias over earlier, but not the other way
// (this way bias=0 replaces, but doesn't forbid later optimizations) // (this way bias=0 replaces, but doesn't forbid later optimizations)
let new_len = (new.image.idat_data.len() as f64 let new_len = (new.image.idat_data.len() as f64
@ -168,9 +78,131 @@ impl Evaluator {
true true
}; };
if is_best { if is_best {
best_result = if new.is_reduction { Some(new) } else { None }; self.best_result = if new.is_reduction { Some(new) } else { None };
} }
} }
best_result.map(|res| res.image)
fn get_result(self) -> Option<PngData> {
self.best_result.map(|res| res.image)
}
}
/// Collect image versions and pick one that compresses best
pub(crate) struct Evaluator {
deadline: Arc<Deadline>,
nth: AtomicUsize,
best_candidate_size: Arc<AtomicMin>,
/// images are sent to the thread for evaluation
#[cfg(feature = "parallel")]
eval_send: SyncSender<Candidate>,
// the thread helps evaluate images asynchronously
#[cfg(feature = "parallel")]
eval_thread: thread::JoinHandle<Option<PngData>>,
// in non-parallel mode, images are evaluated synchronously
#[cfg(not(feature = "parallel"))]
eval_comparator: std::cell::RefCell<Comparator>,
}
impl Evaluator {
pub fn new(deadline: Arc<Deadline>) -> Self {
#[cfg(feature = "parallel")]
let (tx, rx) = sync_channel(4);
Self {
deadline,
best_candidate_size: Arc::new(AtomicMin::new(None)),
nth: AtomicUsize::new(0),
#[cfg(feature = "parallel")]
eval_send: tx,
#[cfg(feature = "parallel")]
eval_thread: thread::spawn(move || {
let mut comparator = Comparator::default();
for candidate in rx {
comparator.evaluate(candidate);
}
comparator.get_result()
}),
#[cfg(not(feature = "parallel"))]
eval_comparator: Default::default(),
}
}
/// Wait for all evaluations to finish and return smallest reduction
/// Or `None` if all reductions were worse than baseline.
#[cfg(feature = "parallel")]
pub fn get_result(self) -> Option<PngData> {
drop(self.eval_send); // disconnect the sender, breaking the loop in the thread
self.eval_thread.join().expect("eval thread")
}
#[cfg(not(feature = "parallel"))]
pub fn get_result(self) -> Option<PngData> {
self.eval_comparator.into_inner().get_result()
}
/// Set baseline image. It will be used only to measure minimum compression level required
pub fn set_baseline(&self, image: Arc<PngImage>) {
self.try_image_inner(image, 1.0, false)
}
/// Check if the image is smaller than others
/// Bias is a value in 0..=1 range. Compressed size is multiplied by
/// this fraction when comparing to the best, so 0.95 allows 5% larger size.
pub fn try_image(&self, image: Arc<PngImage>, bias: f32) {
self.try_image_inner(image, bias, true)
}
fn try_image_inner(&self, image: Arc<PngImage>, bias: f32, is_reduction: bool) {
let nth = self.nth.fetch_add(1, SeqCst);
// These clones are only cheap refcounts
let deadline = self.deadline.clone();
let best_candidate_size = self.best_candidate_size.clone();
// sends it off asynchronously for compression,
// but results will be collected via the message queue
#[cfg(feature = "parallel")]
let eval_send = self.eval_send.clone();
rayon::spawn(move || {
let filters_iter = STD_FILTERS.par_iter().with_max_len(1);
// Updating of best result inside the parallel loop would require locks,
// which are dangerous to do in side Rayon's loop.
// Instead, only update (atomic) best size in real time,
// and the best result later without need for locks.
filters_iter.for_each(|&filter| {
if deadline.passed() {
return;
}
if let Ok(idat_data) = deflate::deflate(
&image.filter_image(filter),
STD_COMPRESSION,
STD_STRATEGY,
STD_WINDOW,
&best_candidate_size,
&deadline,
) {
best_candidate_size.set_min(idat_data.len());
// the rest is shipped to the evavluation/collection thread
let new = Candidate {
image: PngData {
idat_data,
raw: Arc::clone(&image),
},
bias,
filter,
is_reduction,
nth,
};
#[cfg(feature = "parallel")]
{
eval_send.send(new).expect("send");
}
#[cfg(not(feature = "parallel"))]
{
self.eval_comparator.borrow_mut().evaluate(new);
}
}
});
});
} }
} }

View file

@ -779,19 +779,25 @@ fn perform_reductions(
try_alpha_reductions(png, &opts.alphas, eval); try_alpha_reductions(png, &opts.alphas, eval);
} }
struct DeadlineImp {
start: Instant,
timeout: Duration,
print_message: AtomicBool,
}
/// Keep track of processing timeout /// Keep track of processing timeout
pub(crate) struct Deadline { pub(crate) struct Deadline {
start: Instant, imp: Option<DeadlineImp>,
timeout: Option<Duration>,
print_message: AtomicBool,
} }
impl Deadline { impl Deadline {
pub fn new(timeout: Option<Duration>, verbose: bool) -> Self { pub fn new(timeout: Option<Duration>, verbose: bool) -> Self {
Self { Self {
imp: timeout.map(|timeout| DeadlineImp {
start: Instant::now(), start: Instant::now(),
timeout, timeout,
print_message: AtomicBool::new(verbose), print_message: AtomicBool::new(verbose),
})
} }
} }
@ -799,11 +805,11 @@ impl Deadline {
/// ///
/// If the verbose option is on, it also prints a timeout message once. /// If the verbose option is on, it also prints a timeout message once.
pub fn passed(&self) -> bool { pub fn passed(&self) -> bool {
if let Some(timeout) = self.timeout { if let Some(imp) = &self.imp {
let elapsed = self.start.elapsed(); let elapsed = imp.start.elapsed();
if elapsed > timeout { if elapsed > imp.timeout {
if self.print_message.load(Ordering::Relaxed) { if imp.print_message.load(Ordering::Relaxed) {
self.print_message.store(false, Ordering::Relaxed); imp.print_message.store(false, Ordering::Relaxed);
eprintln!("Timed out after {} second(s)", elapsed.as_secs()); eprintln!("Timed out after {} second(s)", elapsed.as_secs());
} }
return true; return true;