use std::{
future::Future,
num::NonZeroUsize,
pin::Pin,
sync::OnceLock,
task::{Context, Poll},
time::Duration,
};
use crate::runtime::{TaskHandle, current_runtime_handle};
#[cfg(not(target_arch = "wasm32"))]
mod native;
#[derive(Clone, Copy, Debug, Eq, PartialEq, thiserror::Error)]
pub enum BlockingError {
#[error("the blocking executor is at capacity")]
Saturated,
#[error("the blocking executor is shut down")]
Shutdown,
#[error("a blocking worker could not start")]
WorkerUnavailable,
#[error("blocking work panicked")]
Panicked,
#[error("blocking work is unsupported on this target")]
Unsupported,
#[error("a blocking UI callback needs a live runtime")]
NoRuntime,
}
#[derive(Clone, Copy, Debug)]
pub struct BlockingExecutorConfig {
pub max_threads: NonZeroUsize,
pub max_queued: NonZeroUsize,
pub idle_timeout: Duration,
}
impl Default for BlockingExecutorConfig {
fn default() -> Self {
let max_threads = std::thread::available_parallelism().unwrap_or(NonZeroUsize::MIN);
Self {
max_threads,
max_queued: max_threads.saturating_mul(NonZeroUsize::new(4).expect("nonzero")),
idle_timeout: Duration::from_secs(30),
}
}
}
#[derive(Clone)]
pub struct BlockingExecutor {
#[cfg(not(target_arch = "wasm32"))]
inner: std::sync::Arc<native::Executor>,
}
impl Default for BlockingExecutor {
fn default() -> Self {
Self::new(BlockingExecutorConfig::default())
}
}
impl BlockingExecutor {
pub fn new(config: BlockingExecutorConfig) -> Self {
#[cfg(not(target_arch = "wasm32"))]
{
Self {
inner: std::sync::Arc::new(native::Executor::new(config)),
}
}
#[cfg(target_arch = "wasm32")]
{
let _ = config;
Self {}
}
}
pub fn try_submit<T, F>(&self, work: F) -> Result<BlockingTask<T>, BlockingError>
where
T: Send + 'static,
F: FnOnce() -> T + Send + 'static,
{
#[cfg(not(target_arch = "wasm32"))]
{
self.inner
.try_submit(work)
.map(|inner| BlockingTask { inner })
}
#[cfg(target_arch = "wasm32")]
{
let _ = work;
Err(BlockingError::Unsupported)
}
}
pub async fn submit<T, F>(&self, work: F) -> Result<T, BlockingError>
where
T: Send + 'static,
F: FnOnce() -> T + Send + 'static,
{
#[cfg(not(target_arch = "wasm32"))]
{
let inner = self.inner.submit(work).await?;
BlockingTask { inner }.await
}
#[cfg(target_arch = "wasm32")]
{
let _ = work;
Err(BlockingError::Unsupported)
}
}
pub fn shutdown(&self) {
#[cfg(not(target_arch = "wasm32"))]
self.inner.shutdown();
}
}
#[must_use = "dropping the task cancels waiting work"]
pub struct BlockingTask<T> {
#[cfg(not(target_arch = "wasm32"))]
inner: native::Pending<T>,
#[cfg(target_arch = "wasm32")]
marker: std::marker::PhantomData<fn() -> T>,
}
impl<T> Future for BlockingTask<T> {
type Output = Result<T, BlockingError>;
fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
#[cfg(not(target_arch = "wasm32"))]
{
self.get_mut().inner.poll(context)
}
#[cfg(target_arch = "wasm32")]
{
let _ = (self, context);
Poll::Ready(Err(BlockingError::Unsupported))
}
}
}
fn shared_executor() -> &'static BlockingExecutor {
static EXECUTOR: OnceLock<BlockingExecutor> = OnceLock::new();
EXECUTOR.get_or_init(BlockingExecutor::default)
}
#[expect(non_snake_case)]
pub async fn withBlocking<T, F>(work: F) -> Result<T, BlockingError>
where
T: Send + 'static,
F: FnOnce() -> T + Send + 'static,
{
shared_executor().submit(work).await
}
#[expect(non_snake_case)]
pub fn launchBlocking<T>(
work: impl FnOnce() -> T + Send + 'static,
on_ui: impl FnOnce(Result<T, BlockingError>) + 'static,
) -> Result<TaskHandle, BlockingError>
where
T: Send + 'static,
{
let runtime = current_runtime_handle().ok_or(BlockingError::NoRuntime)?;
runtime
.spawn_ui(async move { on_ui(withBlocking(work).await) })
.ok_or(BlockingError::NoRuntime)
}