use std::future::Future;
use std::sync::OnceLock;
use tokio::runtime::Runtime;
static RUNTIME: OnceLock<Runtime> = OnceLock::new();
pub fn runtime() -> &'static Runtime {
RUNTIME.get_or_init(|| {
let threads = worker_thread_count(
std::env::var("SAIL_RUNTIME_THREADS").ok().as_deref(),
std::thread::available_parallelism().map_or(2, usize::from),
);
tokio::runtime::Builder::new_multi_thread()
.worker_threads(threads)
.thread_name("sail-core")
.enable_all()
.build()
.expect("failed to build sail-core tokio runtime")
})
}
fn worker_thread_count(env_value: Option<&str>, available: usize) -> usize {
if let Some(raw) = env_value {
if let Ok(n) = raw.trim().parse::<usize>() {
if (1..=256).contains(&n) {
return n;
}
}
}
available.clamp(2, 8)
}
pub fn block_on<F: Future>(future: F) -> F::Output {
debug_assert!(
tokio::runtime::Handle::try_current().is_err(),
"sail::block_on called from within a tokio runtime; use the async API instead"
);
runtime().block_on(future)
}
#[cfg(test)]
mod tests {
use super::worker_thread_count;
#[test]
fn defaults_scale_with_the_machine_between_2_and_8() {
assert_eq!(worker_thread_count(None, 1), 2);
assert_eq!(worker_thread_count(None, 2), 2);
assert_eq!(worker_thread_count(None, 4), 4);
assert_eq!(worker_thread_count(None, 8), 8);
assert_eq!(worker_thread_count(None, 64), 8);
}
#[test]
fn env_override_wins_within_its_accepted_range() {
assert_eq!(worker_thread_count(Some("1"), 64), 1);
assert_eq!(worker_thread_count(Some("16"), 4), 16);
assert_eq!(worker_thread_count(Some(" 32 "), 4), 32);
assert_eq!(worker_thread_count(Some("256"), 4), 256);
}
#[test]
fn bad_env_values_fall_back_to_the_default() {
assert_eq!(worker_thread_count(Some(""), 4), 4);
assert_eq!(worker_thread_count(Some("0"), 4), 4);
assert_eq!(worker_thread_count(Some("257"), 4), 4);
assert_eq!(worker_thread_count(Some("-2"), 4), 4);
assert_eq!(worker_thread_count(Some("two"), 4), 4);
assert_eq!(worker_thread_count(Some("2.5"), 4), 4);
}
}