use super::*;
pub(super) fn compare_sort_keys(a: &Value, b: &Value) -> std::cmp::Ordering {
use std::cmp::Ordering;
match (a, b) {
(Value::Int(x), Value::Int(y)) => x.cmp(y),
(Value::Float(x), Value::Float(y)) => x.partial_cmp(y).unwrap_or(Ordering::Equal),
(Value::Str(x), Value::Str(y)) => x.cmp(y),
_ => Ordering::Equal,
}
}
pub(super) fn par_max_concurrency() -> usize {
let from_env = std::env::var("LEX_PAR_MAX_CONCURRENCY")
.ok()
.and_then(|s| s.parse::<usize>().ok())
.filter(|n| *n > 0);
let default = std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(4);
from_env.unwrap_or(default).min(64)
}
pub(super) fn par_map_run<'a>(
program: &'a Program,
closure: Value,
items: Vec<Value>,
worker_handlers: Vec<Box<dyn EffectHandler + Send>>,
step_limit: u64,
) -> Result<Vec<Value>, VmError> {
if items.is_empty() {
return Ok(Vec::new());
}
let n_workers = worker_handlers.len().min(items.len()).max(1);
let mut buckets: Vec<Vec<(usize, Value)>> = (0..n_workers).map(|_| Vec::new()).collect();
for (i, v) in items.into_iter().enumerate() {
buckets[i % n_workers].push((i, v));
}
let n_total: usize = buckets.iter().map(|b| b.len()).sum();
let results: std::sync::Mutex<Vec<Option<Result<Value, String>>>> =
std::sync::Mutex::new((0..n_total).map(|_| None).collect());
let mut worker_handlers = worker_handlers;
worker_handlers.truncate(n_workers);
type Pair = (Vec<(usize, Value)>, Box<dyn EffectHandler + Send>);
let pairs: Vec<Pair> = buckets.into_iter().zip(worker_handlers).collect();
std::thread::scope(|s| {
let mut handles = Vec::with_capacity(pairs.len());
for (bucket, handler) in pairs {
let closure = closure.clone();
let results = &results;
handles.push(s.spawn(move || {
let handler_for_vm: Box<dyn EffectHandler + 'a> = handler;
let mut vm = Vm::with_handler(program, handler_for_vm);
vm.set_step_limit(step_limit);
for (idx, item) in bucket {
let r = vm
.invoke_closure_value(closure.clone(), vec![item])
.map_err(|e| format!("{e:?}"));
results.lock().unwrap()[idx] = Some(r);
}
}));
}
for h in handles {
h.join().map_err(|_| ()).ok();
}
});
let mut out = Vec::with_capacity(n_total);
let inner = results.into_inner().unwrap();
for r in inner {
match r {
Some(Ok(v)) => out.push(v),
Some(Err(e)) => return Err(VmError::Effect(format!("par_map worker: {e}"))),
None => return Err(VmError::Panic("par_map worker did not produce a result".into())),
}
}
Ok(out)
}