use std::sync::OnceLock;
use rayon::ThreadPool;
fn thread_count() -> usize {
if let Ok(v) = std::env::var("RUST_HDF5_IO_THREADS") {
if let Ok(n) = v.trim().parse::<usize>() {
if n > 0 {
return n;
}
}
}
std::thread::available_parallelism()
.map(|c| c.get() / 2)
.unwrap_or(1)
.max(1)
}
pub(crate) fn io_pool() -> Option<&'static ThreadPool> {
static POOL: OnceLock<Option<ThreadPool>> = OnceLock::new();
POOL.get_or_init(|| {
rayon::ThreadPoolBuilder::new()
.num_threads(thread_count())
.thread_name(|i| format!("hdf5-io-{i}"))
.build()
.ok()
})
.as_ref()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn pool_is_capped_and_reused() {
let n = thread_count();
assert!(n >= 1, "pool must have at least one thread");
if std::env::var_os("RUST_HDF5_IO_THREADS").is_none() {
let cores = std::thread::available_parallelism()
.map(|c| c.get())
.unwrap_or(1);
assert!(
n <= (cores / 2).max(1),
"default pool must be <= half the cores"
);
}
if let Some(pool) = io_pool() {
let a = io_pool().unwrap() as *const ThreadPool;
let b = io_pool().unwrap() as *const ThreadPool;
assert_eq!(a, b, "io_pool must return the same shared pool");
assert_eq!(pool.current_num_threads(), n);
}
}
}