1#[cfg(not(target_arch = "wasm32"))]
12use std::sync::{Condvar, Mutex, PoisonError};
13use std::{
14 cell::{Cell, RefCell},
15 future::Future,
16 pin::Pin,
17 rc::Rc,
18 sync::{
19 Arc, OnceLock,
20 atomic::{AtomicBool, Ordering},
21 },
22 task::{Context, Poll, Waker},
23 time::Duration,
24};
25
26#[cfg(target_arch = "wasm32")]
27use wasm_bindgen::JsCast;
28use web_time::Instant;
29
30use crate::{
31 hooks::{mutableStateOfNeverEqual, remember},
32 runtime::{RuntimeHandle, TaskHandle, current_runtime_handle},
33 state::{MutableState, State},
34};
35
36pub fn spawn_ui_task(future: impl Future<Output = ()> + 'static) -> Option<TaskHandle> {
42 current_runtime_handle().and_then(|runtime| runtime.spawn_ui(future))
43}
44
45#[derive(Clone)]
51pub struct CoroutineScope {
52 inner: Rc<ScopeInner>,
53}
54
55struct ScopeInner {
56 runtime: Option<RuntimeHandle>,
57 tasks: RefCell<Vec<TaskHandle>>,
58 closed: Cell<bool>,
59}
60
61impl Drop for ScopeInner {
62 fn drop(&mut self) {
63 for task in self.tasks.get_mut().drain(..) {
64 task.cancel();
65 }
66 }
67}
68
69struct CompositionScopeOwner(CoroutineScope);
70
71impl Drop for CompositionScopeOwner {
72 fn drop(&mut self) {
73 self.0.inner.closed.set(true);
74 self.0.cancel();
75 }
76}
77
78impl CoroutineScope {
79 pub fn launch(&self, future: impl Future<Output = ()> + 'static) {
82 if self.inner.closed.get() {
83 return;
84 }
85 let Some(runtime) = self.inner.runtime.clone() else {
86 log::warn!("cranpose: a coroutine scope with no runtime dropped its work");
87 return;
88 };
89 self.inner
90 .tasks
91 .borrow_mut()
92 .retain(|task| !task.is_finished());
93 if let Some(handle) = runtime.spawn_ui(future) {
94 self.inner.tasks.borrow_mut().push(handle);
95 }
96 }
97
98 pub fn cancel(&self) {
100 let tasks = std::mem::take(&mut *self.inner.tasks.borrow_mut());
101 for task in tasks {
102 task.cancel();
103 }
104 }
105
106 #[cfg(test)]
107 pub(crate) fn probe_identity(&self) -> usize {
108 Rc::as_ptr(&self.inner) as *const () as usize
109 }
110}
111
112#[expect(non_snake_case)]
114#[track_caller]
115pub fn rememberCoroutineScope() -> CoroutineScope {
116 remember(|| {
117 CompositionScopeOwner(CoroutineScope {
118 inner: Rc::new(ScopeInner {
119 runtime: current_runtime_handle(),
120 tasks: RefCell::new(Vec::new()),
121 closed: Cell::new(false),
122 }),
123 })
124 })
125 .with(|owner| owner.0.clone())
126}
127
128pub fn delay(duration: Duration) -> Delay {
134 Delay {
135 deadline: Instant::now() + duration,
136 armed: false,
137 fired: Arc::new(AtomicBool::new(false)),
138 }
139}
140
141pub struct Delay {
143 deadline: Instant,
144 armed: bool,
145 fired: Arc<AtomicBool>,
146}
147
148impl Future for Delay {
149 type Output = ();
150
151 fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<()> {
152 if self.fired.load(Ordering::Acquire) || Instant::now() >= self.deadline {
153 return Poll::Ready(());
154 }
155 let this = self.get_mut();
156 if !this.armed {
157 this.armed = true;
158 timer().arm(
159 this.deadline,
160 context.waker().clone(),
161 Arc::clone(&this.fired),
162 );
163 }
164 Poll::Pending
165 }
166}
167
168pub async fn interval(period: Duration, mut tick: impl FnMut()) {
173 loop {
174 delay(period).await;
175 tick();
176 }
177}
178
179#[cfg(not(target_arch = "wasm32"))]
180struct Alarm {
181 deadline: Instant,
182 waker: Waker,
183 fired: Arc<AtomicBool>,
184}
185
186struct Timer {
187 #[cfg(not(target_arch = "wasm32"))]
188 alarms: Mutex<Vec<Alarm>>,
189 #[cfg(not(target_arch = "wasm32"))]
190 wake: Condvar,
191}
192
193fn timer() -> &'static Timer {
194 static TIMER: OnceLock<&'static Timer> = OnceLock::new();
195 TIMER.get_or_init(|| {
196 let timer: &'static Timer = Box::leak(Box::new(Timer::new()));
197 timer.start();
198 timer
199 })
200}
201
202#[cfg(not(target_arch = "wasm32"))]
203impl Timer {
204 fn new() -> Self {
205 Self {
206 alarms: Mutex::new(Vec::new()),
207 wake: Condvar::new(),
208 }
209 }
210
211 fn start(&'static self) {
212 std::thread::Builder::new()
213 .name("cranpose-timer".to_string())
214 .spawn(move || self.run())
215 .expect("the timer thread starts");
216 }
217
218 fn run(&self) {
219 let mut alarms = self.alarms.lock().unwrap_or_else(PoisonError::into_inner);
220 loop {
221 let now = Instant::now();
222 let mut due = Vec::new();
223 let mut next: Option<Duration> = None;
224 alarms.retain(|alarm| {
225 if alarm.deadline <= now {
226 due.push((alarm.waker.clone(), Arc::clone(&alarm.fired)));
227 false
228 } else {
229 let remaining = alarm.deadline - now;
230 next = Some(next.map_or(remaining, |current| current.min(remaining)));
231 true
232 }
233 });
234
235 if !due.is_empty() {
236 drop(alarms);
237 for (waker, fired) in due {
238 fired.store(true, Ordering::Release);
239 waker.wake();
240 }
241 alarms = self.alarms.lock().unwrap_or_else(PoisonError::into_inner);
242 continue;
243 }
244
245 alarms = match next {
246 Some(timeout) => {
247 self.wake
248 .wait_timeout(alarms, timeout)
249 .unwrap_or_else(PoisonError::into_inner)
250 .0
251 }
252 None => self
253 .wake
254 .wait(alarms)
255 .unwrap_or_else(PoisonError::into_inner),
256 };
257 }
258 }
259
260 fn arm(&self, deadline: Instant, waker: Waker, fired: Arc<AtomicBool>) {
261 let mut alarms = self.alarms.lock().unwrap_or_else(PoisonError::into_inner);
262 alarms.push(Alarm {
263 deadline,
264 waker,
265 fired,
266 });
267 self.wake.notify_one();
268 }
269}
270
271#[cfg(target_arch = "wasm32")]
272impl Timer {
273 fn new() -> Self {
274 Self {}
275 }
276
277 fn start(&'static self) {}
278
279 fn arm(&self, deadline: Instant, waker: Waker, fired: Arc<AtomicBool>) {
280 let millis = deadline
281 .saturating_duration_since(Instant::now())
282 .as_millis()
283 .min(i32::MAX as u128) as i32;
284 let callback = wasm_bindgen::closure::Closure::once_into_js(move || {
285 fired.store(true, Ordering::Release);
286 waker.wake();
287 });
288 let scheduled = web_sys::window().and_then(|window| {
289 window
290 .set_timeout_with_callback_and_timeout_and_arguments_0(
291 callback.unchecked_ref(),
292 millis,
293 )
294 .ok()
295 });
296 if scheduled.is_none() {
297 log::warn!("cranpose: no window timer is available; the delay resolves immediately");
298 }
299 }
300}
301
302pub struct EventChannel<T: 'static> {
309 shared: Rc<ChannelShared<T>>,
310}
311
312struct ChannelShared<T: 'static> {
313 ready: RefCell<std::collections::VecDeque<T>>,
314 closed: std::cell::Cell<bool>,
315 delivered: std::cell::Cell<usize>,
316 wakers: RefCell<Vec<Waker>>,
317}
318
319impl<T: 'static> ChannelShared<T> {
320 fn wake_all(&self) {
321 for waker in self.wakers.borrow_mut().drain(..) {
322 waker.wake();
323 }
324 }
325}
326
327impl<T: 'static> Default for EventChannel<T> {
328 fn default() -> Self {
329 Self::new()
330 }
331}
332
333impl<T: 'static> EventChannel<T> {
334 pub fn new() -> Self {
336 Self {
337 shared: Rc::new(ChannelShared {
338 ready: RefCell::new(std::collections::VecDeque::new()),
339 closed: std::cell::Cell::new(false),
340 delivered: std::cell::Cell::new(0),
341 wakers: RefCell::new(Vec::new()),
342 }),
343 }
344 }
345
346 pub fn stream(&self) -> EventStream<T> {
348 EventStream {
349 shared: Rc::clone(&self.shared),
350 }
351 }
352
353 pub fn send(&self, event: T) {
355 if self.shared.closed.get() {
356 return;
357 }
358 self.shared.ready.borrow_mut().push_back(event);
359 self.shared.wake_all();
360 }
361
362 pub fn close(&self) {
364 if self.shared.closed.get() {
365 return;
366 }
367 self.shared.closed.set(true);
368 self.shared.wake_all();
369 }
370
371 pub fn is_closed(&self) -> bool {
373 self.shared.closed.get()
374 }
375
376 pub fn pending(&self) -> usize {
378 self.shared.ready.borrow().len()
379 }
380}
381
382pub struct EventStream<T: 'static> {
388 shared: Rc<ChannelShared<T>>,
389}
390
391impl<T: 'static> Clone for EventStream<T> {
392 fn clone(&self) -> Self {
393 Self {
394 shared: Rc::clone(&self.shared),
395 }
396 }
397}
398
399impl<T: 'static> EventStream<T> {
400 pub fn next(&self) -> EventStreamNext<T> {
403 EventStreamNext {
404 shared: Rc::clone(&self.shared),
405 }
406 }
407
408 pub fn delivered(&self) -> usize {
410 self.shared.delivered.get()
411 }
412}
413
414pub struct EventStreamNext<T: 'static> {
416 shared: Rc<ChannelShared<T>>,
417}
418
419impl<T: 'static> Future for EventStreamNext<T> {
420 type Output = Option<T>;
421
422 fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Option<T>> {
423 if let Some(event) = self.shared.ready.borrow_mut().pop_front() {
424 self.shared.delivered.set(self.shared.delivered.get() + 1);
425 return Poll::Ready(Some(event));
426 }
427 if self.shared.closed.get() {
428 return Poll::Ready(None);
429 }
430 self.shared
431 .wakers
432 .borrow_mut()
433 .push(context.waker().clone());
434 Poll::Pending
435 }
436}
437
438#[expect(non_snake_case)]
444#[track_caller]
445pub fn CollectEvents<T, K>(stream: EventStream<T>, key: K, on_event: impl FnMut(T) + 'static)
446where
447 T: 'static,
448 K: PartialEq + 'static,
449{
450 crate::__launched_effect_async_impl(
451 crate::caller_location_key(),
452 std::panic::Location::caller().into(),
453 key,
454 move |_scope| {
455 let mut on_event = on_event;
456 Box::pin(async move {
457 while let Some(event) = stream.next().await {
458 on_event(event);
459 }
460 })
461 },
462 );
463}
464
465#[expect(non_snake_case)]
470#[track_caller]
471pub fn collectAsState<T, K>(stream: EventStream<T>, key: K, initial: T) -> State<T>
472where
473 T: Clone + 'static,
474 K: PartialEq + 'static,
475{
476 let state = remember(|| mutableStateOfNeverEqual(initial)).with(|state| *state);
477 let sink = state;
478 CollectEvents(stream, key, move |event| sink.set(event));
479 state.as_state()
480}
481
482pub struct EventSender<T: Send + 'static> {
490 #[cfg(not(target_arch = "wasm32"))]
491 dispatcher: crate::runtime::UiDispatcher,
492 bridge: u64,
493 _events: std::marker::PhantomData<fn(T)>,
494}
495
496impl<T: Send + 'static> Clone for EventSender<T> {
497 fn clone(&self) -> Self {
498 Self {
499 #[cfg(not(target_arch = "wasm32"))]
500 dispatcher: self.dispatcher.clone(),
501 bridge: self.bridge,
502 _events: std::marker::PhantomData,
503 }
504 }
505}
506
507impl<T: Send + 'static> EventSender<T> {
508 pub fn send(&self, event: T) {
510 let bridge = self.bridge;
511 #[cfg(not(target_arch = "wasm32"))]
512 self.dispatcher
513 .post(move || deliver_bridged::<T>(bridge, event));
514 #[cfg(target_arch = "wasm32")]
515 deliver_bridged::<T>(bridge, event);
516 }
517}
518
519thread_local! {
520 static BRIDGES: RefCell<std::collections::HashMap<u64, Rc<dyn std::any::Any>>> =
521 RefCell::new(std::collections::HashMap::new());
522}
523
524static NEXT_BRIDGE: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1);
525
526fn deliver_bridged<T: Send + 'static>(bridge: u64, event: T) {
527 let channel = BRIDGES.with(|bridges| bridges.borrow().get(&bridge).cloned());
528 let Some(channel) = channel else {
529 log::debug!("event bridge {bridge} is gone, one event dropped");
530 return;
531 };
532 if let Ok(channel) = channel.downcast::<EventChannel<T>>() {
533 channel.send(event);
534 }
535}
536
537struct Bridge<T: Send + 'static> {
538 id: u64,
539 channel: Rc<EventChannel<T>>,
540}
541
542impl<T: Send + 'static> Bridge<T> {
543 fn new() -> Self {
544 let id = NEXT_BRIDGE.fetch_add(1, Ordering::Relaxed);
545 let channel = Rc::new(EventChannel::<T>::new());
546 BRIDGES.with(|bridges| {
547 bridges
548 .borrow_mut()
549 .insert(id, Rc::clone(&channel) as Rc<dyn std::any::Any>)
550 });
551 Self { id, channel }
552 }
553}
554
555impl<T: Send + 'static> Drop for Bridge<T> {
556 fn drop(&mut self) {
557 BRIDGES.with(|bridges| bridges.borrow_mut().remove(&self.id));
558 self.channel.close();
559 }
560}
561
562#[expect(non_snake_case)]
570#[track_caller]
571pub fn rememberEventStream<T, K, R, S>(key: K, subscribe: S) -> EventStream<T>
572where
573 T: Send + 'static,
574 K: PartialEq + 'static,
575 R: 'static,
576 S: FnOnce(EventSender<T>) -> R + 'static,
577{
578 let bridge = remember(Bridge::<T>::new);
579 let (id, stream) = bridge.with(|bridge| (bridge.id, bridge.channel.stream()));
580 #[cfg(not(target_arch = "wasm32"))]
581 let dispatcher = current_runtime_handle().map(|runtime| runtime.dispatcher());
582
583 crate::__disposable_effect_impl(crate::caller_location_key(), key, move |scope| {
584 #[cfg(not(target_arch = "wasm32"))]
585 let Some(dispatcher) = dispatcher else {
586 log::warn!("cranpose: an event stream was remembered without a runtime");
587 return scope.on_dispose(|| {});
588 };
589 let registration = subscribe(EventSender {
590 #[cfg(not(target_arch = "wasm32"))]
591 dispatcher,
592 bridge: id,
593 _events: std::marker::PhantomData,
594 });
595 scope.on_dispose(move || drop(registration))
596 });
597
598 stream
599}
600
601#[expect(non_snake_case)]
608pub async fn withBlocking<T, F>(work: F) -> T
609where
610 T: Send + 'static,
611 F: FnOnce() -> T + Send + 'static,
612{
613 #[cfg(not(target_arch = "wasm32"))]
614 {
615 let slot: Arc<Mutex<Option<T>>> = Arc::new(Mutex::new(None));
616 let done = Arc::new(AtomicBool::new(false));
617 let wakers: Arc<Mutex<Vec<Waker>>> = Arc::new(Mutex::new(Vec::new()));
618
619 let worker_slot = Arc::clone(&slot);
620 let worker_done = Arc::clone(&done);
621 let worker_wakers = Arc::clone(&wakers);
622 BlockingPool::get().submit(Box::new(move || {
623 let value = work();
624 *worker_slot.lock().unwrap_or_else(PoisonError::into_inner) = Some(value);
625 worker_done.store(true, Ordering::Release);
626 for waker in worker_wakers
627 .lock()
628 .unwrap_or_else(PoisonError::into_inner)
629 .drain(..)
630 {
631 waker.wake();
632 }
633 }));
634
635 BlockingWork { slot, done, wakers }.await
636 }
637 #[cfg(target_arch = "wasm32")]
638 {
639 work()
640 }
641}
642
643#[expect(non_snake_case)]
662pub fn launchBlocking<T>(work: impl FnOnce() -> T + Send + 'static, on_ui: impl FnOnce(T) + 'static)
663where
664 T: Send + 'static,
665{
666 let Some(runtime) = current_runtime_handle() else {
667 on_ui(work());
668 return;
669 };
670 let Some(continuation) = runtime.register_ui_cont(on_ui) else {
671 return;
672 };
673 let dispatcher = runtime.dispatcher();
674 #[cfg(not(target_arch = "wasm32"))]
675 BlockingPool::get().submit(Box::new(move || {
676 dispatcher.post_invoke(continuation, work());
677 }));
678 #[cfg(target_arch = "wasm32")]
679 dispatcher.post_invoke(continuation, work());
680}
681
682#[cfg(not(target_arch = "wasm32"))]
683struct BlockingPool {
684 sender: std::sync::mpsc::Sender<BlockingJob>,
685 receiver: Arc<Mutex<std::sync::mpsc::Receiver<BlockingJob>>>,
686 state: Arc<Mutex<PoolState>>,
687}
688
689#[cfg(not(target_arch = "wasm32"))]
690#[derive(Clone, Copy, Default)]
691struct PoolState {
692 alive: usize,
693 outstanding: usize,
694}
695
696#[cfg(not(target_arch = "wasm32"))]
697type BlockingJob = Box<dyn FnOnce() + Send + 'static>;
698
699#[cfg(not(target_arch = "wasm32"))]
700const MAX_BLOCKING_WORKERS: usize = 64;
701
702#[cfg(not(target_arch = "wasm32"))]
703const _: () = assert!(MAX_BLOCKING_WORKERS > 0 && MAX_BLOCKING_WORKERS <= 256);
704
705#[cfg(not(target_arch = "wasm32"))]
706impl BlockingPool {
707 fn get() -> &'static BlockingPool {
708 static POOL: OnceLock<BlockingPool> = OnceLock::new();
709 POOL.get_or_init(BlockingPool::new)
710 }
711
712 fn new() -> BlockingPool {
713 let (sender, receiver) = std::sync::mpsc::channel();
714 BlockingPool {
715 sender,
716 receiver: Arc::new(Mutex::new(receiver)),
717 state: Arc::new(Mutex::new(PoolState::default())),
718 }
719 }
720
721 fn submit(&self, job: BlockingJob) {
722 if self.take_slot() {
723 self.start_worker();
724 }
725 if let Err(returned) = self.sender.send(job) {
726 self.release_slot();
727 (returned.0)();
728 }
729 }
730
731 fn take_slot(&self) -> bool {
732 let mut state = self.state.lock().unwrap_or_else(PoisonError::into_inner);
733 state.outstanding += 1;
734 let grow = state.alive < state.outstanding && state.alive < MAX_BLOCKING_WORKERS;
735 if grow {
736 state.alive += 1;
737 }
738 grow
739 }
740
741 fn release_slot(&self) {
742 let mut state = self.state.lock().unwrap_or_else(PoisonError::into_inner);
743 state.outstanding = state.outstanding.saturating_sub(1);
744 }
745
746 fn start_worker(&self) {
747 let receiver = Arc::clone(&self.receiver);
748 let counters = Arc::clone(&self.state);
749 let started = std::thread::Builder::new()
750 .name("cranpose-blocking".to_string())
751 .spawn(move || {
752 loop {
753 let job = {
754 let queue = receiver.lock().unwrap_or_else(PoisonError::into_inner);
755 queue.recv()
756 };
757 let Ok(job) = job else {
758 break;
759 };
760 job();
761 let mut counters = counters.lock().unwrap_or_else(PoisonError::into_inner);
762 counters.outstanding = counters.outstanding.saturating_sub(1);
763 }
764 });
765 if started.is_err() {
766 let mut state = self.state.lock().unwrap_or_else(PoisonError::into_inner);
767 state.alive -= 1;
768 }
769 }
770}
771
772#[cfg(not(target_arch = "wasm32"))]
773struct BlockingWork<T> {
774 slot: Arc<Mutex<Option<T>>>,
775 done: Arc<AtomicBool>,
776 wakers: Arc<Mutex<Vec<Waker>>>,
777}
778
779#[cfg(not(target_arch = "wasm32"))]
780impl<T> Future for BlockingWork<T> {
781 type Output = T;
782
783 fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<T> {
784 if self.done.load(Ordering::Acquire)
785 && let Some(value) = self
786 .slot
787 .lock()
788 .unwrap_or_else(PoisonError::into_inner)
789 .take()
790 {
791 return Poll::Ready(value);
792 }
793 self.wakers
794 .lock()
795 .unwrap_or_else(PoisonError::into_inner)
796 .push(context.waker().clone());
797 if self.done.load(Ordering::Acquire)
798 && let Some(value) = self
799 .slot
800 .lock()
801 .unwrap_or_else(PoisonError::into_inner)
802 .take()
803 {
804 return Poll::Ready(value);
805 }
806 Poll::Pending
807 }
808}
809
810#[expect(non_snake_case)]
816#[track_caller]
817pub fn produceState<T, K, F>(initial: T, key: K, producer: F) -> State<T>
818where
819 T: Clone + 'static,
820 K: PartialEq + 'static,
821 F: FnOnce(ProduceScope<T>) -> Pin<Box<dyn Future<Output = ()>>> + 'static,
822{
823 let state = remember(|| mutableStateOfNeverEqual(initial)).with(|state| *state);
824 let handle = ProduceScope { state };
825 crate::__launched_effect_async_impl(
826 crate::caller_location_key(),
827 std::panic::Location::caller().into(),
828 key,
829 move |_scope| producer(handle),
830 );
831 state.as_state()
832}
833
834pub struct ProduceScope<T: Clone + 'static> {
836 state: MutableState<T>,
837}
838
839impl<T: Clone + 'static> ProduceScope<T> {
840 pub fn set(&self, value: T) {
842 self.state.set(value);
843 }
844}
845
846#[cfg(test)]
847#[path = "tests/concurrency_tests.rs"]
848mod tests;
849
850#[cfg(test)]
851#[path = "tests/concurrency_stream_tests.rs"]
852mod stream_tests;
853
854#[cfg(test)]
855#[path = "tests/concurrency_timer_race_tests.rs"]
856mod timer_race_tests;