Skip to main content

cranpose_core/
blocking.rs

1use 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/// Failure to admit or complete blocking work.
16#[derive(Clone, Copy, Debug, Eq, PartialEq, thiserror::Error)]
17pub enum BlockingError {
18    /// An immediate `try_submit` found a full queue. The closure has not run.
19    #[error("the blocking executor is at capacity")]
20    Saturated,
21    /// The executor stopped accepting work and discarded its waiting jobs.
22    #[error("the blocking executor is shut down")]
23    Shutdown,
24    /// No worker could be started. The closure has not run.
25    #[error("a blocking worker could not start")]
26    WorkerUnavailable,
27    /// The closure unwound. With panic=abort, a panic still ends the process.
28    #[error("blocking work panicked")]
29    Panicked,
30    /// This target has no supported executor for blocking closures.
31    #[error("blocking work is unsupported on this target")]
32    Unsupported,
33    /// A UI callback was requested without a live composition runtime.
34    #[error("a blocking UI callback needs a live runtime")]
35    NoRuntime,
36}
37
38/// Bounds for a native blocking executor. Running closures cannot be interrupted.
39#[derive(Clone, Copy, Debug)]
40pub struct BlockingExecutorConfig {
41    /// Maximum concurrent worker threads, created only when work needs them.
42    pub max_threads: NonZeroUsize,
43    /// Maximum waiting jobs, in addition to jobs already running.
44    pub max_queued: NonZeroUsize,
45    /// How long an idle worker remains available. Zero retires it immediately.
46    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/// A bounded, lazily started executor for synchronous work.
61///
62/// Clones share one pool. Dropping the last executor or calling `shutdown`
63/// rejects new submissions and fails waiting jobs; running jobs may finish.
64/// The default limits use the available CPU count and four waiting jobs per CPU.
65/// Web submissions return `Unsupported` without executing the closure.
66#[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    /// Creates an executor without starting threads or reserving queue storage.
80    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    /// Admits work without waiting for capacity, or returns `Saturated`.
95    ///
96    /// The returned future owns the job: dropping it removes waiting work and
97    /// releases its captures. A job already taken by a worker may finish.
98    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    /// Waits asynchronously for queue capacity, then runs work and returns its result.
117    ///
118    /// A full queue suspends this future without blocking its polling thread or
119    /// returning `Saturated`. Dropping the future cancels admission or queued work.
120    /// Each suspended caller retains its closure until admission or cancellation;
121    /// applications should bound the number of concurrent producer tasks.
122    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    /// Stops admission and resolves waiting jobs with `Shutdown`.
140    /// Does not wait for or interrupt running closures.
141    pub fn shutdown(&self) {
142        #[cfg(not(target_arch = "wasm32"))]
143        self.inner.shutdown();
144    }
145}
146
147/// Completion of an admitted job. Dropping it cancels work still waiting to run.
148#[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/// Runs synchronous work on the shared bounded executor when first polled.
178///
179/// A full queue suspends this future until capacity is available, keeping the UI
180/// thread free. Shutdown and unwinding panics are reported as errors. Dropping
181/// the future releases its captures and cancels work not yet started. Running
182/// closures cannot be interrupted. Bound concurrent callers to bound the memory
183/// retained by their suspended futures.
184/// On Web this returns `Unsupported` without running the closure.
185#[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/// Starts blocking work and delivers its result on the current runtime's UI
195/// thread. A full queue waits asynchronously inside the runtime-owned task.
196/// Execution errors are delivered to `on_ui`; only a missing runtime is returned
197/// immediately without invoking the callback.
198///
199/// The runtime owns the returned task until completion. Calling its `cancel`
200/// method, or dropping the runtime, cancels waiting work and the callback.
201/// For screen ownership, prefer `rememberCoroutineScope().launch` with
202/// `withBlocking`. Without a live runtime this returns `NoRuntime`.
203#[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}