#![allow(
clippy::missing_const_for_thread_local,
reason = "Melinoe 0.9.0's pinned thread_cached! expansion owns this initializer"
)]
use super::core::IoReactor;
melinoe::thread_cached! {
pub(crate) mod active_reactor: *const IoReactor;
}
pub(crate) static GLOBAL_REACTOR: std::sync::OnceLock<Option<std::sync::Arc<IoReactor>>> =
std::sync::OnceLock::new();
#[cfg(test)]
thread_local! {
static FORCE_NO_REACTOR: std::cell::Cell<bool> = const { std::cell::Cell::new(false) };
}
impl IoReactor {
pub fn with_active<F, R>(&self, f: F) -> R
where
F: FnOnce() -> R,
{
struct Restore(Option<*const IoReactor>);
impl Drop for Restore {
fn drop(&mut self) {
match self.0 {
Some(old) => active_reactor::set(old),
None => active_reactor::clear(),
}
}
}
let _restore = Restore(active_reactor::get());
active_reactor::set(self as *const IoReactor);
f()
}
pub fn get_active() -> Option<&'static IoReactor> {
if let Some(ptr) = active_reactor::get() {
return Some(unsafe { &*ptr });
}
#[cfg(test)]
if FORCE_NO_REACTOR.with(std::cell::Cell::get) {
return None;
}
GLOBAL_REACTOR
.get_or_init(|| {
let reactor = std::sync::Arc::new(IoReactor::new().ok()?);
let driver = std::sync::Arc::clone(&reactor);
std::thread::Builder::new()
.name("moirai-global-reactor".to_string())
.spawn(move || {
let _ = driver.run();
})
.ok()?;
Some(reactor)
})
.as_deref()
}
#[cfg(test)]
pub(crate) fn with_reactor_disabled<F, R>(f: F) -> R
where
F: FnOnce() -> R,
{
struct Restore(bool);
impl Drop for Restore {
fn drop(&mut self) {
FORCE_NO_REACTOR.with(|cell| cell.set(self.0));
}
}
let _restore = Restore(FORCE_NO_REACTOR.with(|cell| cell.replace(true)));
f()
}
}