cranpose_core/
blocking.rs1use std::{
2 future::Future,
3 num::NonZeroUsize,
4 pin::Pin,
5 sync::OnceLock,
6 task::{Context, Poll},
7 time::Duration,
8};
9
10use crate::runtime::{TaskHandle, current_runtime_handle};
11
12#[cfg(not(target_arch = "wasm32"))]
13mod native;
14
15#[derive(Clone, Copy, Debug, Eq, PartialEq, thiserror::Error)]
17pub enum BlockingError {
18 #[error("the blocking executor is at capacity")]
20 Saturated,
21 #[error("the blocking executor is shut down")]
23 Shutdown,
24 #[error("a blocking worker could not start")]
26 WorkerUnavailable,
27 #[error("blocking work panicked")]
29 Panicked,
30 #[error("blocking work is unsupported on this target")]
32 Unsupported,
33 #[error("a blocking UI callback needs a live runtime")]
35 NoRuntime,
36}
37
38#[derive(Clone, Copy, Debug)]
40pub struct BlockingExecutorConfig {
41 pub max_threads: NonZeroUsize,
43 pub max_queued: NonZeroUsize,
45 pub idle_timeout: Duration,
47}
48
49impl Default for BlockingExecutorConfig {
50 fn default() -> Self {
51 let max_threads = std::thread::available_parallelism().unwrap_or(NonZeroUsize::MIN);
52 Self {
53 max_threads,
54 max_queued: max_threads.saturating_mul(NonZeroUsize::new(4).expect("nonzero")),
55 idle_timeout: Duration::from_secs(30),
56 }
57 }
58}
59
60#[derive(Clone)]
67pub struct BlockingExecutor {
68 #[cfg(not(target_arch = "wasm32"))]
69 inner: std::sync::Arc<native::Executor>,
70}
71
72impl Default for BlockingExecutor {
73 fn default() -> Self {
74 Self::new(BlockingExecutorConfig::default())
75 }
76}
77
78impl BlockingExecutor {
79 pub fn new(config: BlockingExecutorConfig) -> Self {
81 #[cfg(not(target_arch = "wasm32"))]
82 {
83 Self {
84 inner: std::sync::Arc::new(native::Executor::new(config)),
85 }
86 }
87 #[cfg(target_arch = "wasm32")]
88 {
89 let _ = config;
90 Self {}
91 }
92 }
93
94 pub fn try_submit<T, F>(&self, work: F) -> Result<BlockingTask<T>, BlockingError>
99 where
100 T: Send + 'static,
101 F: FnOnce() -> T + Send + 'static,
102 {
103 #[cfg(not(target_arch = "wasm32"))]
104 {
105 self.inner
106 .try_submit(work)
107 .map(|inner| BlockingTask { inner })
108 }
109 #[cfg(target_arch = "wasm32")]
110 {
111 let _ = work;
112 Err(BlockingError::Unsupported)
113 }
114 }
115
116 pub async fn submit<T, F>(&self, work: F) -> Result<T, BlockingError>
123 where
124 T: Send + 'static,
125 F: FnOnce() -> T + Send + 'static,
126 {
127 #[cfg(not(target_arch = "wasm32"))]
128 {
129 let inner = self.inner.submit(work).await?;
130 BlockingTask { inner }.await
131 }
132 #[cfg(target_arch = "wasm32")]
133 {
134 let _ = work;
135 Err(BlockingError::Unsupported)
136 }
137 }
138
139 pub fn shutdown(&self) {
142 #[cfg(not(target_arch = "wasm32"))]
143 self.inner.shutdown();
144 }
145}
146
147#[must_use = "dropping the task cancels waiting work"]
149pub struct BlockingTask<T> {
150 #[cfg(not(target_arch = "wasm32"))]
151 inner: native::Pending<T>,
152 #[cfg(target_arch = "wasm32")]
153 marker: std::marker::PhantomData<fn() -> T>,
154}
155
156impl<T> Future for BlockingTask<T> {
157 type Output = Result<T, BlockingError>;
158
159 fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
160 #[cfg(not(target_arch = "wasm32"))]
161 {
162 self.get_mut().inner.poll(context)
163 }
164 #[cfg(target_arch = "wasm32")]
165 {
166 let _ = (self, context);
167 Poll::Ready(Err(BlockingError::Unsupported))
168 }
169 }
170}
171
172fn shared_executor() -> &'static BlockingExecutor {
173 static EXECUTOR: OnceLock<BlockingExecutor> = OnceLock::new();
174 EXECUTOR.get_or_init(BlockingExecutor::default)
175}
176
177#[expect(non_snake_case)]
186pub async fn withBlocking<T, F>(work: F) -> Result<T, BlockingError>
187where
188 T: Send + 'static,
189 F: FnOnce() -> T + Send + 'static,
190{
191 shared_executor().submit(work).await
192}
193
194#[expect(non_snake_case)]
204pub fn launchBlocking<T>(
205 work: impl FnOnce() -> T + Send + 'static,
206 on_ui: impl FnOnce(Result<T, BlockingError>) + 'static,
207) -> Result<TaskHandle, BlockingError>
208where
209 T: Send + 'static,
210{
211 let runtime = current_runtime_handle().ok_or(BlockingError::NoRuntime)?;
212 runtime
213 .spawn_ui(async move { on_ui(withBlocking(work).await) })
214 .ok_or(BlockingError::NoRuntime)
215}