Skip to main content

cranpose_core/
concurrency.rs

1//! Composition-scoped structured concurrency.
2//!
3//! `LaunchedEffect` covers work that starts because a key changed. This covers
4//! the rest: work started from an event handler, timed work, work that feeds a
5//! piece of state, and blocking work that must not run on the UI thread.
6//! Everything here is owned by the composition — a scope cancels its tasks when
7//! it leaves, and a timer stops when nothing is waiting on it — so an
8//! application never keeps its own task list or its own "is this still alive"
9//! flag.
10
11#[cfg(not(any(target_arch = "wasm32", target_vendor = "apple")))]
12use std::sync::Condvar;
13#[cfg(not(target_arch = "wasm32"))]
14use std::sync::{Mutex, PoisonError};
15use std::{
16    cell::{Cell, RefCell},
17    future::Future,
18    pin::Pin,
19    rc::Rc,
20    sync::{
21        Arc, OnceLock,
22        atomic::{AtomicBool, Ordering},
23    },
24    task::{Context, Poll, Waker},
25    time::Duration,
26};
27
28#[cfg(target_arch = "wasm32")]
29use wasm_bindgen::JsCast;
30use web_time::Instant;
31
32use crate::{
33    hooks::{mutableStateOfNeverEqual, remember},
34    runtime::{RuntimeHandle, TaskHandle, current_runtime_handle},
35    state::{MutableState, State},
36};
37
38/// Spawns `future` on the current runtime's UI task queue.
39///
40/// Framework-internal: application code launches through a
41/// [`CoroutineScope`] so the work is cancelled with its
42/// composition. Returns `None` when there is no runtime on this thread.
43pub fn spawn_ui_task(future: impl Future<Output = ()> + 'static) -> Option<TaskHandle> {
44    current_runtime_handle().and_then(|runtime| runtime.spawn_ui(future))
45}
46
47/// A cancellation scope for work launched outside the composition pass.
48///
49/// Tasks launched through a scope are cancelled when the scope leaves the
50/// composition, so a click handler can start an asynchronous job without the
51/// job outliving the screen that started it.
52#[derive(Clone)]
53pub struct CoroutineScope {
54    inner: Rc<ScopeInner>,
55}
56
57struct ScopeInner {
58    runtime: Option<RuntimeHandle>,
59    tasks: RefCell<Vec<TaskHandle>>,
60    closed: Cell<bool>,
61}
62
63impl Drop for ScopeInner {
64    fn drop(&mut self) {
65        for task in self.tasks.get_mut().drain(..) {
66            task.cancel();
67        }
68    }
69}
70
71struct CompositionScopeOwner(CoroutineScope);
72
73impl Drop for CompositionScopeOwner {
74    fn drop(&mut self) {
75        self.0.inner.closed.set(true);
76        self.0.cancel();
77    }
78}
79
80impl CoroutineScope {
81    /// Launches `future`, keeping it alive until it finishes or the scope is
82    /// cancelled.
83    pub fn launch(&self, future: impl Future<Output = ()> + 'static) {
84        if self.inner.closed.get() {
85            return;
86        }
87        let Some(runtime) = self.inner.runtime.clone() else {
88            log::warn!("cranpose: a coroutine scope with no runtime dropped its work");
89            return;
90        };
91        self.inner
92            .tasks
93            .borrow_mut()
94            .retain(|task| !task.is_finished());
95        if let Some(handle) = runtime.spawn_ui(future) {
96            self.inner.tasks.borrow_mut().push(handle);
97        }
98    }
99
100    /// Cancels every task this scope launched.
101    pub fn cancel(&self) {
102        let tasks = std::mem::take(&mut *self.inner.tasks.borrow_mut());
103        for task in tasks {
104            task.cancel();
105        }
106    }
107
108    #[cfg(test)]
109    pub(crate) fn probe_identity(&self) -> usize {
110        Rc::as_ptr(&self.inner) as *const () as usize
111    }
112}
113
114/// Remembers a [`CoroutineScope`] bound to this position in the composition.
115#[expect(non_snake_case)]
116#[track_caller]
117pub fn rememberCoroutineScope() -> CoroutineScope {
118    remember(|| {
119        CompositionScopeOwner(CoroutineScope {
120            inner: Rc::new(ScopeInner {
121                runtime: current_runtime_handle(),
122                tasks: RefCell::new(Vec::new()),
123                closed: Cell::new(false),
124            }),
125        })
126    })
127    .with(|owner| owner.0.clone())
128}
129
130/// Resolves after `duration` has elapsed.
131///
132/// The wait is served by the framework's timer, which posts the wake-up onto
133/// the runtime's UI queue. Nothing spins and no frames are requested while a
134/// delay is pending, so a one-minute timer costs nothing for a minute.
135pub fn delay(duration: Duration) -> Delay {
136    Delay {
137        deadline: Instant::now() + duration,
138        armed: false,
139        fired: Arc::new(AtomicBool::new(false)),
140    }
141}
142
143/// The future returned by [`delay`].
144pub struct Delay {
145    deadline: Instant,
146    armed: bool,
147    fired: Arc<AtomicBool>,
148}
149
150impl Future for Delay {
151    type Output = ();
152
153    fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<()> {
154        if self.fired.load(Ordering::Acquire) || Instant::now() >= self.deadline {
155            return Poll::Ready(());
156        }
157        let this = self.get_mut();
158        if !this.armed {
159            this.armed = true;
160            timer().arm(
161                this.deadline,
162                context.waker().clone(),
163                Arc::clone(&this.fired),
164            );
165        }
166        Poll::Pending
167    }
168}
169
170/// Runs `tick` every `period` until the returned future is dropped.
171///
172/// The first tick happens after one full period, matching a repeating timer
173/// rather than a leading-edge one.
174pub async fn interval(period: Duration, mut tick: impl FnMut()) {
175    loop {
176        delay(period).await;
177        tick();
178    }
179}
180
181#[cfg(not(any(target_arch = "wasm32", target_vendor = "apple")))]
182struct Alarm {
183    deadline: Instant,
184    waker: Waker,
185    fired: Arc<AtomicBool>,
186}
187
188struct Timer {
189    #[cfg(not(any(target_arch = "wasm32", target_vendor = "apple")))]
190    alarms: Mutex<Vec<Alarm>>,
191    #[cfg(not(any(target_arch = "wasm32", target_vendor = "apple")))]
192    wake: Condvar,
193}
194
195fn timer() -> &'static Timer {
196    static TIMER: OnceLock<&'static Timer> = OnceLock::new();
197    TIMER.get_or_init(|| {
198        let timer: &'static Timer = Box::leak(Box::new(Timer::new()));
199        timer.start();
200        timer
201    })
202}
203
204#[cfg(not(any(target_arch = "wasm32", target_vendor = "apple")))]
205impl Timer {
206    fn new() -> Self {
207        Self {
208            alarms: Mutex::new(Vec::new()),
209            wake: Condvar::new(),
210        }
211    }
212
213    fn start(&'static self) {
214        std::thread::Builder::new()
215            .name("cranpose-timer".to_string())
216            .spawn(move || self.run())
217            .expect("the timer thread starts");
218    }
219
220    fn run(&self) {
221        let mut alarms = self.alarms.lock().unwrap_or_else(PoisonError::into_inner);
222        loop {
223            let now = Instant::now();
224            let mut due = Vec::new();
225            let mut next: Option<Duration> = None;
226            alarms.retain(|alarm| {
227                if alarm.deadline <= now {
228                    due.push((alarm.waker.clone(), Arc::clone(&alarm.fired)));
229                    false
230                } else {
231                    let remaining = alarm.deadline - now;
232                    next = Some(next.map_or(remaining, |current| current.min(remaining)));
233                    true
234                }
235            });
236
237            if !due.is_empty() {
238                drop(alarms);
239                for (waker, fired) in due {
240                    fired.store(true, Ordering::Release);
241                    waker.wake();
242                }
243                alarms = self.alarms.lock().unwrap_or_else(PoisonError::into_inner);
244                continue;
245            }
246
247            alarms = match next {
248                Some(timeout) => {
249                    self.wake
250                        .wait_timeout(alarms, timeout)
251                        .unwrap_or_else(PoisonError::into_inner)
252                        .0
253                }
254                None => self
255                    .wake
256                    .wait(alarms)
257                    .unwrap_or_else(PoisonError::into_inner),
258            };
259        }
260    }
261
262    fn arm(&self, deadline: Instant, waker: Waker, fired: Arc<AtomicBool>) {
263        let mut alarms = self.alarms.lock().unwrap_or_else(PoisonError::into_inner);
264        alarms.push(Alarm {
265            deadline,
266            waker,
267            fired,
268        });
269        self.wake.notify_one();
270    }
271}
272
273/// Apple platforms hand each deadline to Grand Central Dispatch on the
274/// user-interactive queue. A thread sleeping on a condition variable at the
275/// default quality of service wakes about 8 ms late there, as macOS coalesces
276/// its timers, so a one-frame `delay` often missed its frame; GCD at
277/// user-interactive quality of service wakes within a millisecond.
278#[cfg(target_vendor = "apple")]
279impl Timer {
280    fn new() -> Self {
281        Self {}
282    }
283
284    fn start(&'static self) {}
285
286    fn arm(&self, deadline: Instant, waker: Waker, fired: Arc<AtomicBool>) {
287        use dispatch2::{DispatchQoS, DispatchQueue, DispatchTime, GlobalQueueIdentifier};
288
289        let nanos = deadline
290            .saturating_duration_since(Instant::now())
291            .as_nanos()
292            .min(i64::MAX as u128) as i64;
293        let queue = DispatchQueue::global_queue(GlobalQueueIdentifier::QualityOfService(
294            DispatchQoS::UserInteractive,
295        ));
296        let fire = move || {
297            fired.store(true, Ordering::Release);
298            waker.wake();
299        };
300        if queue.after(DispatchTime::NOW.time(nanos), fire).is_err() {
301            log::error!("cranpose: GCD refused a timer; the delay never resolves");
302        }
303    }
304}
305
306#[cfg(target_arch = "wasm32")]
307impl Timer {
308    fn new() -> Self {
309        Self {}
310    }
311
312    fn start(&'static self) {}
313
314    fn arm(&self, deadline: Instant, waker: Waker, fired: Arc<AtomicBool>) {
315        let millis = deadline
316            .saturating_duration_since(Instant::now())
317            .as_millis()
318            .min(i32::MAX as u128) as i32;
319        let callback = wasm_bindgen::closure::Closure::once_into_js(move || {
320            fired.store(true, Ordering::Release);
321            waker.wake();
322        });
323        let scheduled = web_sys::window().and_then(|window| {
324            window
325                .set_timeout_with_callback_and_timeout_and_arguments_0(
326                    callback.unchecked_ref(),
327                    millis,
328                )
329                .ok()
330        });
331        if scheduled.is_none() {
332            log::warn!("cranpose: no window timer is available; the delay resolves immediately");
333        }
334    }
335}
336
337/// The producing half of an [`EventStream`].
338///
339/// A service that receives events from outside the composition — a platform
340/// callback, a worker thread, a socket — publishes through a channel, and every
341/// pending collector is woken. This is the shape that replaces
342/// "register an observer, then drain a queue" everywhere in the framework.
343pub struct EventChannel<T: 'static> {
344    shared: Rc<ChannelShared<T>>,
345}
346
347struct ChannelShared<T: 'static> {
348    ready: RefCell<std::collections::VecDeque<T>>,
349    closed: std::cell::Cell<bool>,
350    delivered: std::cell::Cell<usize>,
351    wakers: RefCell<Vec<Waker>>,
352}
353
354impl<T: 'static> ChannelShared<T> {
355    fn wake_all(&self) {
356        for waker in self.wakers.borrow_mut().drain(..) {
357            waker.wake();
358        }
359    }
360}
361
362impl<T: 'static> Default for EventChannel<T> {
363    fn default() -> Self {
364        Self::new()
365    }
366}
367
368impl<T: 'static> EventChannel<T> {
369    /// Creates an open channel.
370    pub fn new() -> Self {
371        Self {
372            shared: Rc::new(ChannelShared {
373                ready: RefCell::new(std::collections::VecDeque::new()),
374                closed: std::cell::Cell::new(false),
375                delivered: std::cell::Cell::new(0),
376                wakers: RefCell::new(Vec::new()),
377            }),
378        }
379    }
380
381    /// The consuming half, handed to collectors.
382    pub fn stream(&self) -> EventStream<T> {
383        EventStream {
384            shared: Rc::clone(&self.shared),
385        }
386    }
387
388    /// Publishes one event and wakes every pending collector.
389    pub fn send(&self, event: T) {
390        if self.shared.closed.get() {
391            return;
392        }
393        self.shared.ready.borrow_mut().push_back(event);
394        self.shared.wake_all();
395    }
396
397    /// Ends the stream. Collectors drain what is queued and then finish.
398    pub fn close(&self) {
399        if self.shared.closed.get() {
400            return;
401        }
402        self.shared.closed.set(true);
403        self.shared.wake_all();
404    }
405
406    /// Whether the channel has been closed.
407    pub fn is_closed(&self) -> bool {
408        self.shared.closed.get()
409    }
410
411    /// How many events are queued but not yet taken.
412    pub fn pending(&self) -> usize {
413        self.shared.ready.borrow().len()
414    }
415}
416
417/// The consuming half of an [`EventChannel`].
418///
419/// Collectors take events one at a time; an event goes to exactly one
420/// collector, so two collectors share the stream rather than each seeing every
421/// event.
422pub struct EventStream<T: 'static> {
423    shared: Rc<ChannelShared<T>>,
424}
425
426impl<T: 'static> Clone for EventStream<T> {
427    fn clone(&self) -> Self {
428        Self {
429            shared: Rc::clone(&self.shared),
430        }
431    }
432}
433
434impl<T: 'static> EventStream<T> {
435    /// Resolves with the next event, or `None` once the stream is closed and
436    /// drained.
437    pub fn next(&self) -> EventStreamNext<T> {
438        EventStreamNext {
439            shared: Rc::clone(&self.shared),
440        }
441    }
442
443    /// How many events this stream has handed out.
444    pub fn delivered(&self) -> usize {
445        self.shared.delivered.get()
446    }
447}
448
449/// The future returned by [`EventStream::next`].
450pub struct EventStreamNext<T: 'static> {
451    shared: Rc<ChannelShared<T>>,
452}
453
454impl<T: 'static> Future for EventStreamNext<T> {
455    type Output = Option<T>;
456
457    fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Option<T>> {
458        if let Some(event) = self.shared.ready.borrow_mut().pop_front() {
459            self.shared.delivered.set(self.shared.delivered.get() + 1);
460            return Poll::Ready(Some(event));
461        }
462        if self.shared.closed.get() {
463            return Poll::Ready(None);
464        }
465        self.shared
466            .wakers
467            .borrow_mut()
468            .push(context.waker().clone());
469        Poll::Pending
470    }
471}
472
473/// Collects `stream` for as long as this call stays in the composition,
474/// handing each event to `on_event`.
475///
476/// `key` re-starts the collection when it changes, exactly like
477/// `LaunchedEffect`.
478#[expect(non_snake_case)]
479#[track_caller]
480pub fn CollectEvents<T, K>(stream: EventStream<T>, key: K, on_event: impl FnMut(T) + 'static)
481where
482    T: 'static,
483    K: PartialEq + 'static,
484{
485    crate::__launched_effect_async_impl(
486        crate::caller_location_key(),
487        std::panic::Location::caller().into(),
488        key,
489        move |_scope| {
490            let mut on_event = on_event;
491            Box::pin(async move {
492                while let Some(event) = stream.next().await {
493                    on_event(event);
494                }
495            })
496        },
497    );
498}
499
500/// Collects `stream` into state, starting at `initial`.
501///
502/// The composition reads the latest value the stream produced, and recomposes
503/// when a new one arrives.
504#[expect(non_snake_case)]
505#[track_caller]
506pub fn collectAsState<T, K>(stream: EventStream<T>, key: K, initial: T) -> State<T>
507where
508    T: Clone + 'static,
509    K: PartialEq + 'static,
510{
511    let state = remember(|| mutableStateOfNeverEqual(initial)).with(|state| *state);
512    let sink = state;
513    CollectEvents(stream, key, move |event| sink.set(event));
514    state.as_state()
515}
516
517/// A `Send` publishing handle for a composition-scoped [`EventStream`].
518///
519/// Platform services publish events from whatever thread they run on — a JNI
520/// callback, a worker, a socket reader. The sender hops each event onto the UI
521/// thread through the runtime's dispatcher and pushes it into the stream the
522/// composition is collecting, so no service and no application ever writes that
523/// hop again.
524pub struct EventSender<T: Send + 'static> {
525    #[cfg(not(target_arch = "wasm32"))]
526    dispatcher: crate::runtime::UiDispatcher,
527    bridge: u64,
528    _events: std::marker::PhantomData<fn(T)>,
529}
530
531impl<T: Send + 'static> Clone for EventSender<T> {
532    fn clone(&self) -> Self {
533        Self {
534            #[cfg(not(target_arch = "wasm32"))]
535            dispatcher: self.dispatcher.clone(),
536            bridge: self.bridge,
537            _events: std::marker::PhantomData,
538        }
539    }
540}
541
542impl<T: Send + 'static> EventSender<T> {
543    /// Publishes `event` to the composition that owns this bridge.
544    pub fn send(&self, event: T) {
545        let bridge = self.bridge;
546        #[cfg(not(target_arch = "wasm32"))]
547        self.dispatcher
548            .post(move || deliver_bridged::<T>(bridge, event));
549        #[cfg(target_arch = "wasm32")]
550        deliver_bridged::<T>(bridge, event);
551    }
552}
553
554thread_local! {
555    static BRIDGES: RefCell<std::collections::HashMap<u64, Rc<dyn std::any::Any>>> =
556        RefCell::new(std::collections::HashMap::new());
557}
558
559static NEXT_BRIDGE: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1);
560
561fn deliver_bridged<T: Send + 'static>(bridge: u64, event: T) {
562    let channel = BRIDGES.with(|bridges| bridges.borrow().get(&bridge).cloned());
563    let Some(channel) = channel else {
564        log::debug!("event bridge {bridge} is gone, one event dropped");
565        return;
566    };
567    if let Ok(channel) = channel.downcast::<EventChannel<T>>() {
568        channel.send(event);
569    }
570}
571
572struct Bridge<T: Send + 'static> {
573    id: u64,
574    channel: Rc<EventChannel<T>>,
575}
576
577impl<T: Send + 'static> Bridge<T> {
578    fn new() -> Self {
579        let id = NEXT_BRIDGE.fetch_add(1, Ordering::Relaxed);
580        let channel = Rc::new(EventChannel::<T>::new());
581        BRIDGES.with(|bridges| {
582            bridges
583                .borrow_mut()
584                .insert(id, Rc::clone(&channel) as Rc<dyn std::any::Any>)
585        });
586        Self { id, channel }
587    }
588}
589
590impl<T: Send + 'static> Drop for Bridge<T> {
591    fn drop(&mut self) {
592        BRIDGES.with(|bridges| bridges.borrow_mut().remove(&self.id));
593        self.channel.close();
594    }
595}
596
597/// Turns a platform subscription into a composition-scoped [`EventStream`].
598///
599/// `subscribe` receives a `Send` [`EventSender`] and returns whatever
600/// registration handle the service uses; that handle is dropped — unsubscribing
601/// the service — when `key` changes or the composition leaves. This is the one
602/// place the framework bridges "a service publishes from another thread" to
603/// "a composition collects".
604#[expect(non_snake_case)]
605#[track_caller]
606pub fn rememberEventStream<T, K, R, S>(key: K, subscribe: S) -> EventStream<T>
607where
608    T: Send + 'static,
609    K: PartialEq + 'static,
610    R: 'static,
611    S: FnOnce(EventSender<T>) -> R + 'static,
612{
613    let bridge = remember(Bridge::<T>::new);
614    let (id, stream) = bridge.with(|bridge| (bridge.id, bridge.channel.stream()));
615    #[cfg(not(target_arch = "wasm32"))]
616    let dispatcher = current_runtime_handle().map(|runtime| runtime.dispatcher());
617
618    crate::__disposable_effect_impl(crate::caller_location_key(), key, move |scope| {
619        #[cfg(not(target_arch = "wasm32"))]
620        let Some(dispatcher) = dispatcher else {
621            log::warn!("cranpose: an event stream was remembered without a runtime");
622            return scope.on_dispose(|| {});
623        };
624        let registration = subscribe(EventSender {
625            #[cfg(not(target_arch = "wasm32"))]
626            dispatcher,
627            bridge: id,
628            _events: std::marker::PhantomData,
629        });
630        scope.on_dispose(move || drop(registration))
631    });
632
633    stream
634}
635
636/// Runs `work` off the UI thread and resolves with its result on the UI thread.
637///
638/// This is the escape hatch for genuinely blocking work — parsing a large file,
639/// a synchronous provider call — that must not stall composition. On the web
640/// there is one thread, so `work` runs inline; callers keep the unit of work
641/// small enough that this is honest on every target.
642#[expect(non_snake_case)]
643pub async fn withBlocking<T, F>(work: F) -> T
644where
645    T: Send + 'static,
646    F: FnOnce() -> T + Send + 'static,
647{
648    #[cfg(not(target_arch = "wasm32"))]
649    {
650        let slot: Arc<Mutex<Option<T>>> = Arc::new(Mutex::new(None));
651        let done = Arc::new(AtomicBool::new(false));
652        let wakers: Arc<Mutex<Vec<Waker>>> = Arc::new(Mutex::new(Vec::new()));
653
654        let worker_slot = Arc::clone(&slot);
655        let worker_done = Arc::clone(&done);
656        let worker_wakers = Arc::clone(&wakers);
657        BlockingPool::get().submit(Box::new(move || {
658            let value = work();
659            *worker_slot.lock().unwrap_or_else(PoisonError::into_inner) = Some(value);
660            worker_done.store(true, Ordering::Release);
661            for waker in worker_wakers
662                .lock()
663                .unwrap_or_else(PoisonError::into_inner)
664                .drain(..)
665            {
666                waker.wake();
667            }
668        }));
669
670        BlockingWork { slot, done, wakers }.await
671    }
672    #[cfg(target_arch = "wasm32")]
673    {
674        work()
675    }
676}
677
678/// Runs `work` off the UI thread and hands its result to `on_ui` on the UI
679/// thread.
680///
681/// The callback shape of [`withBlocking`], for the code that is not already in
682/// a coroutine: an event handler, a button, anything that wants to start some
683/// blocking work and carry on. Both share the same pool, so an application
684/// that uses one, the other, or both never spends more than one set of threads
685/// on blocking work.
686///
687/// Without a runtime — a unit test, a tool — `work` runs inline and `on_ui`
688/// follows it, so a caller behaves the same either way.
689///
690/// ```rust,ignore
691/// launchBlocking(
692///     move || std::fs::read(path),
693///     move |bytes| document.set(bytes.ok()),
694/// );
695/// ```
696#[expect(non_snake_case)]
697pub fn launchBlocking<T>(work: impl FnOnce() -> T + Send + 'static, on_ui: impl FnOnce(T) + 'static)
698where
699    T: Send + 'static,
700{
701    let Some(runtime) = current_runtime_handle() else {
702        on_ui(work());
703        return;
704    };
705    let Some(continuation) = runtime.register_ui_cont(on_ui) else {
706        return;
707    };
708    let dispatcher = runtime.dispatcher();
709    #[cfg(not(target_arch = "wasm32"))]
710    BlockingPool::get().submit(Box::new(move || {
711        dispatcher.post_invoke(continuation, work());
712    }));
713    #[cfg(target_arch = "wasm32")]
714    dispatcher.post_invoke(continuation, work());
715}
716
717#[cfg(not(target_arch = "wasm32"))]
718struct BlockingPool {
719    sender: std::sync::mpsc::Sender<BlockingJob>,
720    receiver: Arc<Mutex<std::sync::mpsc::Receiver<BlockingJob>>>,
721    state: Arc<Mutex<PoolState>>,
722}
723
724#[cfg(not(target_arch = "wasm32"))]
725#[derive(Clone, Copy, Default)]
726struct PoolState {
727    alive: usize,
728    outstanding: usize,
729}
730
731#[cfg(not(target_arch = "wasm32"))]
732type BlockingJob = Box<dyn FnOnce() + Send + 'static>;
733
734#[cfg(not(target_arch = "wasm32"))]
735const MAX_BLOCKING_WORKERS: usize = 64;
736
737#[cfg(not(target_arch = "wasm32"))]
738const _: () = assert!(MAX_BLOCKING_WORKERS > 0 && MAX_BLOCKING_WORKERS <= 256);
739
740#[cfg(not(target_arch = "wasm32"))]
741impl BlockingPool {
742    fn get() -> &'static BlockingPool {
743        static POOL: OnceLock<BlockingPool> = OnceLock::new();
744        POOL.get_or_init(BlockingPool::new)
745    }
746
747    fn new() -> BlockingPool {
748        let (sender, receiver) = std::sync::mpsc::channel();
749        BlockingPool {
750            sender,
751            receiver: Arc::new(Mutex::new(receiver)),
752            state: Arc::new(Mutex::new(PoolState::default())),
753        }
754    }
755
756    fn submit(&self, job: BlockingJob) {
757        if self.take_slot() {
758            self.start_worker();
759        }
760        if let Err(returned) = self.sender.send(job) {
761            self.release_slot();
762            (returned.0)();
763        }
764    }
765
766    fn take_slot(&self) -> bool {
767        let mut state = self.state.lock().unwrap_or_else(PoisonError::into_inner);
768        state.outstanding += 1;
769        let grow = state.alive < state.outstanding && state.alive < MAX_BLOCKING_WORKERS;
770        if grow {
771            state.alive += 1;
772        }
773        grow
774    }
775
776    fn release_slot(&self) {
777        let mut state = self.state.lock().unwrap_or_else(PoisonError::into_inner);
778        state.outstanding = state.outstanding.saturating_sub(1);
779    }
780
781    fn start_worker(&self) {
782        let receiver = Arc::clone(&self.receiver);
783        let counters = Arc::clone(&self.state);
784        let started = std::thread::Builder::new()
785            .name("cranpose-blocking".to_string())
786            .spawn(move || {
787                loop {
788                    let job = {
789                        let queue = receiver.lock().unwrap_or_else(PoisonError::into_inner);
790                        queue.recv()
791                    };
792                    let Ok(job) = job else {
793                        break;
794                    };
795                    job();
796                    let mut counters = counters.lock().unwrap_or_else(PoisonError::into_inner);
797                    counters.outstanding = counters.outstanding.saturating_sub(1);
798                }
799            });
800        if started.is_err() {
801            let mut state = self.state.lock().unwrap_or_else(PoisonError::into_inner);
802            state.alive -= 1;
803        }
804    }
805}
806
807#[cfg(not(target_arch = "wasm32"))]
808struct BlockingWork<T> {
809    slot: Arc<Mutex<Option<T>>>,
810    done: Arc<AtomicBool>,
811    wakers: Arc<Mutex<Vec<Waker>>>,
812}
813
814#[cfg(not(target_arch = "wasm32"))]
815impl<T> Future for BlockingWork<T> {
816    type Output = T;
817
818    fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<T> {
819        if self.done.load(Ordering::Acquire)
820            && let Some(value) = self
821                .slot
822                .lock()
823                .unwrap_or_else(PoisonError::into_inner)
824                .take()
825        {
826            return Poll::Ready(value);
827        }
828        self.wakers
829            .lock()
830            .unwrap_or_else(PoisonError::into_inner)
831            .push(context.waker().clone());
832        if self.done.load(Ordering::Acquire)
833            && let Some(value) = self
834                .slot
835                .lock()
836                .unwrap_or_else(PoisonError::into_inner)
837                .take()
838        {
839            return Poll::Ready(value);
840        }
841        Poll::Pending
842    }
843}
844
845/// Runs `producer` when `key` changes and exposes what it publishes as state.
846///
847/// The Compose `produceState` contract: the producer receives a handle it uses
848/// to publish values, and is cancelled when the key changes or the composition
849/// leaves.
850#[expect(non_snake_case)]
851#[track_caller]
852pub fn produceState<T, K, F>(initial: T, key: K, producer: F) -> State<T>
853where
854    T: Clone + 'static,
855    K: PartialEq + 'static,
856    F: FnOnce(ProduceScope<T>) -> Pin<Box<dyn Future<Output = ()>>> + 'static,
857{
858    let state = remember(|| mutableStateOfNeverEqual(initial)).with(|state| *state);
859    let handle = ProduceScope { state };
860    crate::__launched_effect_async_impl(
861        crate::caller_location_key(),
862        std::panic::Location::caller().into(),
863        key,
864        move |_scope| producer(handle),
865    );
866    state.as_state()
867}
868
869/// The publishing half handed to a [`produceState`] producer.
870pub struct ProduceScope<T: Clone + 'static> {
871    state: MutableState<T>,
872}
873
874impl<T: Clone + 'static> ProduceScope<T> {
875    /// Publishes `value` to the produced state.
876    pub fn set(&self, value: T) {
877        self.state.set(value);
878    }
879}
880
881#[cfg(test)]
882#[path = "tests/concurrency_tests.rs"]
883mod tests;
884
885#[cfg(test)]
886#[path = "tests/concurrency_stream_tests.rs"]
887mod stream_tests;
888
889#[cfg(test)]
890#[path = "tests/concurrency_timer_race_tests.rs"]
891mod timer_race_tests;