1#[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
38pub 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#[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 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 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#[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
130pub 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
143pub 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
170pub 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#[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
337pub 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 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 pub fn stream(&self) -> EventStream<T> {
383 EventStream {
384 shared: Rc::clone(&self.shared),
385 }
386 }
387
388 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 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 pub fn is_closed(&self) -> bool {
408 self.shared.closed.get()
409 }
410
411 pub fn pending(&self) -> usize {
413 self.shared.ready.borrow().len()
414 }
415}
416
417pub 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 pub fn next(&self) -> EventStreamNext<T> {
438 EventStreamNext {
439 shared: Rc::clone(&self.shared),
440 }
441 }
442
443 pub fn delivered(&self) -> usize {
445 self.shared.delivered.get()
446 }
447}
448
449pub 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#[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#[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
517pub 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 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#[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#[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#[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#[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
869pub struct ProduceScope<T: Clone + 'static> {
871 state: MutableState<T>,
872}
873
874impl<T: Clone + 'static> ProduceScope<T> {
875 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;