1use std::{
2 any::Any,
3 cell::{Cell, RefCell},
4 collections::VecDeque,
5 future::Future,
6 pin::Pin,
7 rc::{Rc, Weak},
8 sync::{
9 Arc,
10 atomic::{AtomicBool, AtomicUsize, Ordering},
11 mpsc,
12 },
13 task::{Context, Poll, Waker},
14 thread::ThreadId,
15 thread_local,
16};
17
18#[cfg(any(feature = "internal", test))]
19use crate::frame_clock::FrameClock;
20use crate::{
21 Applier, Command, FrameCallbackId, Key, MutableStateInner, NodeError, RecomposeScope,
22 RecomposeScopeInner, ScopeId,
23 collections::map::HashMap,
24 platform::{RuntimeScheduler, SchedulerRef},
25 state::{MutationPolicy, NeverEqual},
26};
27
28#[derive(Clone, Copy, PartialEq, Eq)]
29pub(crate) enum FrameCallbackKind {
30 Transient,
31 Perpetual,
32}
33
34enum UiMessage {
35 Task(Box<dyn FnOnce() + Send + 'static>),
36 Invoke { id: u64, value: Box<dyn Any + Send> },
37}
38
39type UiContinuation = Box<dyn Fn(Box<dyn Any>) -> bool + 'static>;
40type UiContinuationMap = HashMap<u64, UiContinuation>;
41
42struct TypedStateCell<T: Clone + 'static> {
43 inner: MutableStateInner<T>,
44}
45
46trait ScopeWatchCell {
47 fn unregister_scope(&self, scope_id: ScopeId);
48}
49
50impl<T: Clone + 'static> ScopeWatchCell for TypedStateCell<T> {
51 fn unregister_scope(&self, scope_id: ScopeId) {
52 self.inner.unregister_scope(scope_id);
53 }
54}
55
56struct StateArenaSlot {
57 generation: u32,
58 cell: Option<Rc<dyn Any>>,
59 watcher_cell: Option<Rc<dyn ScopeWatchCell>>,
60 lease: Option<Weak<StateHandleLease>>,
61}
62
63#[derive(Default)]
64struct StateArenaInner {
65 cells: Vec<StateArenaSlot>,
66 free: Vec<u32>,
67}
68
69#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
70pub struct StateArenaDebugStats {
71 pub cells_len: usize,
72 pub cells_cap: usize,
73 pub free_len: usize,
74 pub free_cap: usize,
75}
76
77#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
78pub struct RuntimeDebugStats {
79 pub node_updates_len: usize,
80 pub node_updates_cap: usize,
81 pub invalid_scopes_len: usize,
82 pub scope_queue_len: usize,
83 pub scope_queue_cap: usize,
84 pub frame_callbacks_len: usize,
85 pub frame_callbacks_cap: usize,
86 pub local_tasks_len: usize,
87 pub local_tasks_cap: usize,
88 pub ui_conts_len: usize,
89 pub ui_conts_cap: usize,
90 pub tasks_len: usize,
91 pub tasks_cap: usize,
92 pub external_state_owners_len: usize,
93 pub external_state_owners_cap: usize,
94 pub ui_dispatcher_pending: usize,
95}
96
97#[derive(Default)]
98pub(crate) struct StateArena {
99 inner: RefCell<StateArenaInner>,
100}
101
102impl StateArena {
103 pub(crate) fn alloc<T: Clone + 'static>(&self, value: T, runtime: RuntimeHandle) -> StateId {
104 self.alloc_with_policy(value, runtime, Rc::new(NeverEqual))
105 }
106
107 pub(crate) fn alloc_with_policy<T: Clone + 'static>(
108 &self,
109 value: T,
110 runtime: RuntimeHandle,
111 policy: Rc<dyn MutationPolicy<T>>,
112 ) -> StateId {
113 let (slot, generation) = {
114 let mut inner = self.inner.borrow_mut();
115 loop {
116 let Some(slot) = inner.free.pop() else {
117 let slot = inner.cells.len() as u32;
118 inner.cells.push(StateArenaSlot {
119 generation: 0,
120 cell: None,
121 watcher_cell: None,
122 lease: None,
123 });
124 break (slot, 0);
125 };
126
127 let Some(entry) = inner.cells.get_mut(slot as usize) else {
128 continue;
129 };
130 if entry.cell.is_some() {
131 continue;
132 }
133
134 entry.watcher_cell = None;
135 entry.lease = None;
136 entry.generation = entry.generation.wrapping_add(1);
137 break (slot, entry.generation);
138 }
139 };
140 let id = StateId::new(slot, generation);
141 let inner = MutableStateInner::new_with_policy(value, runtime, policy);
142 inner.install_snapshot_observer(id);
143 let typed_cell = Rc::new(TypedStateCell { inner });
144 let cell: Rc<dyn Any> = typed_cell.clone();
145 let watcher_cell: Rc<dyn ScopeWatchCell> = typed_cell;
146 let mut arena = self.inner.borrow_mut();
147 let slot_entry = &mut arena.cells[slot as usize];
148 slot_entry.cell = Some(cell);
149 slot_entry.watcher_cell = Some(watcher_cell);
150 id
151 }
152
153 fn get_cell_opt(&self, id: StateId) -> Option<Rc<dyn Any>> {
154 self.inner
155 .borrow()
156 .cells
157 .get(id.slot_index())
158 .filter(|cell| cell.generation == id.generation())
159 .and_then(|cell| cell.cell.as_ref())
160 .cloned()
161 }
162
163 fn get_typed<T: Clone + 'static>(&self, id: StateId) -> Rc<TypedStateCell<T>> {
164 match self.get_cell_opt(id) {
165 None => panic!(
166 "state cell missing: slot={}, gen={}, expected={}",
167 id.slot(),
168 id.generation(),
169 std::any::type_name::<T>(),
170 ),
171 Some(cell) => Rc::downcast::<TypedStateCell<T>>(cell).unwrap_or_else(|_| {
172 panic!(
173 "state cell type mismatch: slot={}, gen={}, expected={}",
174 id.slot(),
175 id.generation(),
176 std::any::type_name::<T>(),
177 )
178 }),
179 }
180 }
181
182 fn get_typed_opt<T: Clone + 'static>(&self, id: StateId) -> Option<Rc<TypedStateCell<T>>> {
183 Rc::downcast::<TypedStateCell<T>>(self.get_cell_opt(id)?).ok()
184 }
185
186 pub(crate) fn with_typed<T: Clone + 'static, R>(
187 &self,
188 id: StateId,
189 f: impl FnOnce(&MutableStateInner<T>) -> R,
190 ) -> R {
191 let cell = self.get_typed::<T>(id);
192 f(&cell.inner)
193 }
194
195 pub(crate) fn with_typed_opt<T: Clone + 'static, R>(
196 &self,
197 id: StateId,
198 f: impl FnOnce(&MutableStateInner<T>) -> R,
199 ) -> Option<R> {
200 let cell = self.get_typed_opt::<T>(id)?;
201 Some(f(&cell.inner))
202 }
203
204 pub(crate) fn release(&self, id: StateId) {
205 let cell = {
206 let mut inner = self.inner.borrow_mut();
207 let Some(slot) = inner.cells.get_mut(id.slot_index()) else {
208 return;
209 };
210 if slot.generation != id.generation() {
211 return;
212 }
213 slot.lease = None;
214 slot.watcher_cell = None;
215 let cell = slot.cell.take();
216 if cell.is_some() {
217 inner.free.push(id.slot());
218 }
219 cell
220 };
221 drop(cell);
222 }
223
224 pub(crate) fn stats(&self) -> (usize, usize) {
225 let inner = self.inner.borrow();
226 (inner.cells.len(), inner.free.len())
227 }
228
229 pub(crate) fn debug_stats(&self) -> StateArenaDebugStats {
230 let inner = self.inner.borrow();
231 StateArenaDebugStats {
232 cells_len: inner.cells.len(),
233 cells_cap: inner.cells.capacity(),
234 free_len: inner.free.len(),
235 free_cap: inner.free.capacity(),
236 }
237 }
238
239 pub(crate) fn unregister_scope(&self, id: StateId, scope_id: ScopeId) {
240 let watcher_cell = {
241 let inner = self.inner.borrow();
242 inner
243 .cells
244 .get(id.slot_index())
245 .filter(|slot| slot.generation == id.generation())
246 .and_then(|slot| slot.watcher_cell.as_ref())
247 .cloned()
248 };
249 if let Some(watcher_cell) = watcher_cell {
250 watcher_cell.unregister_scope(scope_id);
251 }
252 }
253
254 pub(crate) fn register_lease(&self, id: StateId, lease: &Rc<StateHandleLease>) {
255 let mut inner = self.inner.borrow_mut();
256 let Some(slot) = inner.cells.get_mut(id.slot_index()) else {
257 return;
258 };
259 if slot.generation != id.generation() {
260 return;
261 }
262 slot.lease = Some(Rc::downgrade(lease));
263 }
264
265 pub(crate) fn retain_lease(&self, id: StateId) -> Option<Rc<StateHandleLease>> {
266 let inner = self.inner.borrow();
267 let slot = inner.cells.get(id.slot_index())?;
268 if slot.generation != id.generation() {
269 return None;
270 }
271 slot.lease.as_ref()?.upgrade()
272 }
273}
274
275#[derive(Copy, Clone, Debug, PartialEq, Eq, Hash)]
276pub struct StateId {
277 slot: u32,
278 generation: u32,
279}
280
281impl StateId {
282 const fn new(slot: u32, generation: u32) -> Self {
283 Self { slot, generation }
284 }
285
286 pub(crate) const fn slot(self) -> u32 {
287 self.slot
288 }
289
290 pub(crate) const fn slot_index(self) -> usize {
291 self.slot as usize
292 }
293
294 pub(crate) const fn generation(self) -> u32 {
295 self.generation
296 }
297}
298
299#[derive(Copy, Clone, Debug, PartialEq, Eq, Hash)]
300pub struct RuntimeId(u32);
301
302impl RuntimeId {
303 fn next() -> Self {
304 NEXT_RUNTIME_ID.with(|next| {
305 let id = next.get();
306 next.set(id.wrapping_add(1));
307 Self(id)
308 })
309 }
310}
311
312struct UiDispatcherInner {
313 scheduler: SchedulerRef,
314 tx: mpsc::Sender<UiMessage>,
315 pending: AtomicUsize,
316}
317
318#[cfg(not(target_arch = "wasm32"))]
319type UiDispatcherRef = Arc<UiDispatcherInner>;
320
321#[cfg(target_arch = "wasm32")]
322type UiDispatcherRef = Rc<UiDispatcherInner>;
323
324impl UiDispatcherInner {
325 fn new(scheduler: SchedulerRef, tx: mpsc::Sender<UiMessage>) -> Self {
326 Self {
327 scheduler,
328 tx,
329 pending: AtomicUsize::new(0),
330 }
331 }
332
333 fn post(&self, task: impl FnOnce() + Send + 'static) {
334 self.pending.fetch_add(1, Ordering::SeqCst);
335 if self.tx.send(UiMessage::Task(Box::new(task))).is_ok() {
336 self.scheduler.schedule_frame();
337 } else {
338 self.pending.fetch_sub(1, Ordering::SeqCst);
339 }
340 }
341
342 fn post_invoke(&self, id: u64, value: Box<dyn Any + Send>) {
343 self.pending.fetch_add(1, Ordering::SeqCst);
344 if self.tx.send(UiMessage::Invoke { id, value }).is_ok() {
345 self.scheduler.schedule_frame();
346 } else {
347 self.pending.fetch_sub(1, Ordering::SeqCst);
348 }
349 }
350
351 fn has_pending(&self) -> bool {
352 self.pending.load(Ordering::SeqCst) > 0
353 }
354}
355
356struct PendingGuard<'a> {
357 counter: &'a AtomicUsize,
358}
359
360impl<'a> PendingGuard<'a> {
361 fn new(counter: &'a AtomicUsize) -> Self {
362 Self { counter }
363 }
364}
365
366impl Drop for PendingGuard<'_> {
367 fn drop(&mut self) {
368 let mut current = self.counter.load(Ordering::SeqCst);
369 loop {
370 if current == 0 {
371 return;
372 }
373 match self.counter.compare_exchange(
374 current,
375 current - 1,
376 Ordering::SeqCst,
377 Ordering::SeqCst,
378 ) {
379 Ok(_) => return,
380 Err(next) => current = next,
381 }
382 }
383 }
384}
385
386#[derive(Clone)]
387pub struct UiDispatcher {
388 inner: UiDispatcherRef,
389}
390
391impl UiDispatcher {
392 fn new(inner: UiDispatcherRef) -> Self {
393 Self { inner }
394 }
395
396 pub fn post(&self, task: impl FnOnce() + Send + 'static) {
397 self.inner.post(task);
398 }
399
400 pub fn post_invoke<T>(&self, id: u64, value: T)
401 where
402 T: Send + 'static,
403 {
404 self.inner.post_invoke(id, Box::new(value));
405 }
406
407 pub fn has_pending(&self) -> bool {
408 self.inner.has_pending()
409 }
410}
411
412struct RuntimeInner {
413 scheduler: SchedulerRef,
414 needs_frame: RefCell<bool>,
415 node_updates: RefCell<Vec<Command>>,
416 invalid_scope_count: Cell<usize>,
417 scope_queue: RefCell<Vec<Weak<RecomposeScopeInner>>>,
418 frame_callbacks: RefCell<VecDeque<FrameCallbackEntry>>,
419 next_frame_callback_id: Cell<u64>,
420 last_frame_time_nanos: Cell<Option<u64>>,
421 ui_dispatcher: UiDispatcherRef,
422 ui_rx: RefCell<mpsc::Receiver<UiMessage>>,
423 local_tasks: RefCell<VecDeque<Box<dyn FnOnce() + 'static>>>,
424 ui_conts: RefCell<UiContinuationMap>,
425 next_cont_id: Cell<u64>,
426 ui_thread_id: ThreadId,
427 tasks: RefCell<HashMap<u64, TaskEntry>>,
428 task_order: RefCell<Vec<u64>>,
429 next_task_id: Cell<u64>,
430 state_arena: StateArena,
431 external_state_owners: RefCell<HashMap<StateId, Rc<StateHandleLease>>>,
432 live_recompose_scope_count: Cell<usize>,
433 forgotten_movables: RefCell<Vec<Key>>,
434 next_movable_content_id: Cell<u64>,
435 runtime_id: RuntimeId,
436}
437
438struct TaskEntry {
439 label: String,
440 future: Option<Pin<Box<dyn Future<Output = ()> + 'static>>>,
441 runnable: Arc<AtomicBool>,
442 waker: Waker,
443}
444
445thread_local! {
446 static NEXT_TASK_LABEL: RefCell<Option<String>> = const { RefCell::new(None) };
447}
448
449pub fn label_next_ui_task(label: impl Into<String>) {
450 NEXT_TASK_LABEL.with(|held| *held.borrow_mut() = Some(label.into()));
451}
452
453impl RuntimeInner {
454 fn new(scheduler: SchedulerRef) -> Self {
455 let (tx, rx) = mpsc::channel();
456 let dispatcher = UiDispatcherRef::new(UiDispatcherInner::new(scheduler.clone(), tx));
457 Self {
458 scheduler,
459 needs_frame: RefCell::new(false),
460 node_updates: RefCell::new(Vec::new()),
461 invalid_scope_count: Cell::new(0),
462 scope_queue: RefCell::new(Vec::new()),
463 frame_callbacks: RefCell::new(VecDeque::new()),
464 next_frame_callback_id: Cell::new(1),
465 last_frame_time_nanos: Cell::new(None),
466 ui_dispatcher: dispatcher,
467 ui_rx: RefCell::new(rx),
468 local_tasks: RefCell::new(VecDeque::new()),
469 ui_conts: RefCell::new(UiContinuationMap::default()),
470 next_cont_id: Cell::new(1),
471 ui_thread_id: std::thread::current().id(),
472 tasks: RefCell::new(HashMap::default()),
473 task_order: RefCell::new(Vec::new()),
474 next_task_id: Cell::new(1),
475 state_arena: StateArena::default(),
476 external_state_owners: RefCell::new(HashMap::default()),
477 live_recompose_scope_count: Cell::new(0),
478 forgotten_movables: RefCell::new(Vec::new()),
479 next_movable_content_id: Cell::new(1),
480 runtime_id: RuntimeId::next(),
481 }
482 }
483
484 fn schedule(&self) {
485 *self.needs_frame.borrow_mut() = true;
486 self.scheduler.schedule_frame();
487 }
488
489 fn enqueue_update(&self, command: Command) {
490 self.node_updates.borrow_mut().push(command);
491 self.schedule();
492 }
493
494 fn take_updates(&self) -> Vec<Command> {
495 self.node_updates.borrow_mut().drain(..).collect::<Vec<_>>()
496 }
497
498 fn has_updates(&self) -> bool {
499 !self.node_updates.borrow().is_empty() || self.has_invalid_scopes()
500 }
501
502 fn register_invalid_scope(&self, scope: Weak<RecomposeScopeInner>) {
503 self.invalid_scope_count
504 .set(self.invalid_scope_count.get() + 1);
505 self.scope_queue.borrow_mut().push(scope);
506 self.schedule();
507 }
508
509 fn requeue_invalid_scope(&self, scope: &RecomposeScope) {
510 if scope.is_enqueued() {
511 self.scope_queue.borrow_mut().push(scope.downgrade());
512 self.schedule();
513 }
514 }
515
516 fn mark_scope_recomposed(&self) {
517 self.invalid_scope_count
518 .set(self.invalid_scope_count.get().saturating_sub(1));
519 }
520
521 fn take_invalidated_scopes(&self) -> Option<Vec<RecomposeScope>> {
522 let mut pending = std::mem::take(&mut *self.scope_queue.borrow_mut());
523 if pending.is_empty() {
524 return None;
525 }
526 let scopes: Vec<RecomposeScope> = pending
527 .iter()
528 .filter_map(RecomposeScope::upgrade)
529 .filter(RecomposeScope::is_enqueued)
530 .collect();
531 pending.clear();
532 let mut queue = self.scope_queue.borrow_mut();
533 if queue.is_empty() {
534 std::mem::swap(&mut *queue, &mut pending);
535 }
536 drop(queue);
537 (!scopes.is_empty()).then_some(scopes)
538 }
539
540 fn has_invalid_scopes(&self) -> bool {
541 self.invalid_scope_count.get() != 0
542 }
543
544 fn queued_invalid_scope_ids(&self) -> Vec<ScopeId> {
545 let mut ids: Vec<ScopeId> = self
546 .scope_queue
547 .borrow()
548 .iter()
549 .filter_map(RecomposeScope::upgrade)
550 .filter(RecomposeScope::is_enqueued)
551 .map(|scope| scope.id())
552 .collect();
553 ids.sort_unstable();
554 ids.dedup();
555 ids
556 }
557
558 fn increment_live_recompose_scope_count(&self) {
559 self.live_recompose_scope_count
560 .set(self.live_recompose_scope_count.get().saturating_add(1));
561 }
562
563 fn decrement_live_recompose_scope_count(&self) {
564 self.live_recompose_scope_count
565 .set(self.live_recompose_scope_count.get().saturating_sub(1));
566 }
567
568 fn live_recompose_scope_count(&self) -> usize {
569 self.live_recompose_scope_count.get()
570 }
571
572 fn has_frame_callbacks(&self) -> bool {
573 !self.frame_callbacks.borrow().is_empty()
574 }
575
576 fn has_transient_frame_callbacks(&self) -> bool {
577 self.frame_callbacks
578 .borrow()
579 .iter()
580 .any(|entry| entry.kind == FrameCallbackKind::Transient)
581 }
582
583 fn enqueue_ui_task(&self, task: Box<dyn FnOnce() + 'static>) {
584 self.local_tasks.borrow_mut().push_back(task);
585 self.schedule();
586 }
587
588 fn spawn_ui_task(&self, future: Pin<Box<dyn Future<Output = ()> + 'static>>) -> u64 {
589 let id = self.next_task_id.get();
590 self.next_task_id.set(id + 1);
591 let label = NEXT_TASK_LABEL
592 .with(|held| held.borrow_mut().take())
593 .unwrap_or_else(|| "unnamed".to_string());
594 let runnable = Arc::new(AtomicBool::new(true));
595 let waker = RuntimeTaskWaker::new(self, Arc::clone(&runnable)).into_waker();
596 self.tasks.borrow_mut().insert(
597 id,
598 TaskEntry {
599 label,
600 future: Some(future),
601 runnable,
602 waker,
603 },
604 );
605 self.task_order.borrow_mut().push(id);
606 self.schedule();
607 id
608 }
609
610 fn cancel_task(&self, id: u64) {
611 let task = self.tasks.borrow_mut().remove(&id);
612 self.task_order.borrow_mut().retain(|queued| *queued != id);
613 drop(task);
614 }
615
616 fn has_task(&self, id: u64) -> bool {
617 self.tasks
618 .try_borrow()
619 .map_or(true, |tasks| tasks.contains_key(&id))
620 }
621
622 fn poll_async_tasks(&self) -> bool {
623 let order = std::mem::take(&mut *self.task_order.borrow_mut());
624 let mut pending = Vec::with_capacity(order.len());
625 let mut made_progress = false;
626 for id in order {
627 let task = {
628 let mut tasks = self.tasks.borrow_mut();
629 let Some(entry) = tasks.get_mut(&id) else {
630 continue;
631 };
632 if entry.runnable.swap(false, Ordering::AcqRel) {
633 entry
634 .future
635 .take()
636 .map(|future| (future, entry.waker.clone()))
637 } else {
638 None
639 }
640 };
641 let Some((mut future, waker)) = task else {
642 pending.push(id);
643 continue;
644 };
645 let mut cx = Context::from_waker(&waker);
646 match future.as_mut().poll(&mut cx) {
647 Poll::Ready(()) => {
648 self.cancel_task(id);
649 made_progress = true;
650 }
651 Poll::Pending => {
652 let mut tasks = self.tasks.borrow_mut();
653 if let Some(entry) = tasks.get_mut(&id) {
654 entry.future = Some(future);
655 pending.push(id);
656 } else {
657 drop(tasks);
658 drop(future);
659 }
660 }
661 }
662 }
663 if !pending.is_empty() {
664 pending.retain(|id| self.has_task(*id));
665 self.task_order.borrow_mut().extend(pending);
666 }
667 made_progress
668 }
669
670 fn drain_ui(&self) {
671 loop {
672 let mut executed = false;
673
674 {
675 let rx = &mut *self.ui_rx.borrow_mut();
676 for message in rx.try_iter() {
677 executed = true;
678 let _guard = PendingGuard::new(&self.ui_dispatcher.pending);
679 match message {
680 UiMessage::Task(task) => {
681 task();
682 }
683 UiMessage::Invoke { id, value } => {
684 self.invoke_ui_cont(id, value);
685 }
686 }
687 }
688 }
689
690 if self.run_local_tasks() {
691 executed = true;
692 }
693
694 if self.poll_async_tasks() {
695 executed = true;
696 }
697 self.run_local_tasks();
704
705 if !executed {
706 break;
707 }
708 }
709
710 self.clear_needs_frame_if_idle();
711 }
712
713 fn run_local_tasks(&self) -> bool {
714 let mut ran = false;
715 loop {
716 let task = self.local_tasks.borrow_mut().pop_front();
717 let Some(task) = task else {
718 return ran;
719 };
720 ran = true;
721 task();
722 }
723 }
724
725 fn has_pending_ui(&self) -> bool {
726 let local_pending = self
727 .local_tasks
728 .try_borrow()
729 .map_or(true, |tasks| !tasks.is_empty());
730
731 local_pending || self.ui_dispatcher.has_pending() || self.has_runnable_tasks()
732 }
733
734 fn has_runnable_tasks(&self) -> bool {
735 self.tasks.try_borrow().map_or(true, |tasks| {
736 tasks
737 .values()
738 .any(|task| task.runnable.load(Ordering::Acquire))
739 })
740 }
741
742 fn register_ui_cont<T: 'static>(&self, f: impl FnOnce(T) + 'static) -> u64 {
743 debug_assert_eq!(
744 std::thread::current().id(),
745 self.ui_thread_id,
746 "UI continuation registered off the runtime thread",
747 );
748 let id = self.next_cont_id.get();
749 self.next_cont_id.set(id + 1);
750 let callback = RefCell::new(Some(f));
751 self.ui_conts.borrow_mut().insert(
752 id,
753 Box::new(move |value: Box<dyn Any>| {
754 let Ok(value) = value.downcast::<T>() else {
755 return false;
756 };
757 let Some(slot) = callback.borrow_mut().take() else {
758 return true;
759 };
760 slot(*value);
761 true
762 }),
763 );
764 id
765 }
766
767 fn invoke_ui_cont(&self, id: u64, value: Box<dyn Any + Send>) {
768 debug_assert_eq!(
769 std::thread::current().id(),
770 self.ui_thread_id,
771 "UI continuation invoked off the runtime thread",
772 );
773 let callback = { self.ui_conts.borrow_mut().remove(&id) };
774 if let Some(callback) = callback {
775 let value: Box<dyn Any> = value;
776 if !callback(value) {
777 self.ui_conts.borrow_mut().insert(id, callback);
778 }
779 }
780 }
781
782 fn cancel_ui_cont(&self, id: u64) {
783 let continuation = self.ui_conts.borrow_mut().remove(&id);
784 drop(continuation);
785 }
786
787 fn register_frame_callback(
788 &self,
789 kind: FrameCallbackKind,
790 callback: Box<dyn FnOnce(u64) + 'static>,
791 ) -> FrameCallbackId {
792 let id = self.next_frame_callback_id.get();
793 self.next_frame_callback_id.set(id + 1);
794 self.frame_callbacks
795 .borrow_mut()
796 .push_back(FrameCallbackEntry {
797 id,
798 kind,
799 callback: Some(callback),
800 });
801 self.schedule();
802 id
803 }
804
805 fn cancel_frame_callback(&self, id: FrameCallbackId) {
806 let removed = {
807 let mut callbacks = self.frame_callbacks.borrow_mut();
808 callbacks
809 .iter()
810 .position(|entry| entry.id == id)
811 .and_then(|index| callbacks.remove(index))
812 };
813 drop(removed);
814 self.clear_needs_frame_if_idle();
815 }
816
817 fn clear_needs_frame_if_idle(&self) {
818 if !self.has_invalid_scopes()
819 && !self.has_updates()
820 && !self.has_frame_callbacks()
821 && !self.has_pending_ui()
822 {
823 *self.needs_frame.borrow_mut() = false;
824 }
825 }
826
827 fn drain_frame_callbacks(&self, frame_time_nanos: u64) {
828 if self
829 .last_frame_time_nanos
830 .get()
831 .is_some_and(|previous| frame_time_nanos <= previous)
832 {
833 return;
834 }
835 self.last_frame_time_nanos.set(Some(frame_time_nanos));
836 let next_frame_id = self.next_frame_callback_id.get();
837 if self.has_frame_callbacks() {
838 let _ = crate::run_in_mutable_snapshot(|| {
839 loop {
840 let entry = {
841 let mut callbacks = self.frame_callbacks.borrow_mut();
842 if callbacks
843 .front()
844 .is_some_and(|entry| entry.id < next_frame_id)
845 {
846 callbacks.pop_front()
847 } else {
848 None
849 }
850 };
851 let Some(mut entry) = entry else {
852 break;
853 };
854 if let Some(callback) = entry.callback.take() {
855 callback(frame_time_nanos);
856 }
857 }
858 });
859 }
860
861 self.clear_needs_frame_if_idle();
862 }
863
864 fn debug_stats(&self) -> RuntimeDebugStats {
865 let node_updates = self.node_updates.borrow();
866 let scope_queue = self.scope_queue.borrow();
867 let frame_callbacks = self.frame_callbacks.borrow();
868 let local_tasks = self.local_tasks.borrow();
869 let ui_conts = self.ui_conts.borrow();
870 let tasks = self.tasks.borrow();
871 let external_state_owners = self.external_state_owners.borrow();
872
873 RuntimeDebugStats {
874 node_updates_len: node_updates.len(),
875 node_updates_cap: node_updates.capacity(),
876 invalid_scopes_len: self.invalid_scope_count.get(),
877 scope_queue_len: scope_queue.len(),
878 scope_queue_cap: scope_queue.capacity(),
879 frame_callbacks_len: frame_callbacks.len(),
880 frame_callbacks_cap: frame_callbacks.capacity(),
881 local_tasks_len: local_tasks.len(),
882 local_tasks_cap: local_tasks.capacity(),
883 ui_conts_len: ui_conts.len(),
884 ui_conts_cap: ui_conts.capacity(),
885 tasks_len: tasks.len(),
886 tasks_cap: tasks.capacity(),
887 external_state_owners_len: external_state_owners.len(),
888 external_state_owners_cap: external_state_owners.capacity(),
889 ui_dispatcher_pending: self.ui_dispatcher.pending.load(Ordering::SeqCst),
890 }
891 }
892}
893
894#[derive(Clone)]
895pub struct Runtime {
896 inner: Rc<RuntimeInner>,
897}
898
899impl Runtime {
900 pub fn new(scheduler: SchedulerRef) -> Self {
901 let inner = Rc::new(RuntimeInner::new(scheduler));
902 let runtime = Self { inner };
903 let handle = runtime.handle();
904 register_runtime_handle(&handle);
905 LAST_RUNTIME.with(|slot| *slot.borrow_mut() = Some(handle));
906 runtime
907 }
908
909 pub fn handle(&self) -> RuntimeHandle {
910 RuntimeHandle {
911 inner: Rc::downgrade(&self.inner),
912 dispatcher: UiDispatcher::new(self.inner.ui_dispatcher.clone()),
913 ui_thread_id: self.inner.ui_thread_id,
914 id: self.inner.runtime_id,
915 }
916 }
917
918 pub fn has_updates(&self) -> bool {
919 self.inner.has_updates()
920 }
921
922 pub fn needs_frame(&self) -> bool {
923 *self.inner.needs_frame.borrow() || self.inner.has_runnable_tasks()
924 }
925
926 pub fn set_needs_frame(&self, value: bool) {
927 *self.inner.needs_frame.borrow_mut() = value;
928 }
929
930 pub fn last_frame_time_nanos(&self) -> Option<u64> {
935 self.inner.last_frame_time_nanos.get()
936 }
937
938 #[cfg(any(feature = "internal", test))]
939 pub fn frame_clock(&self) -> FrameClock {
940 FrameClock::new(self.handle())
941 }
942}
943
944impl Drop for Runtime {
945 fn drop(&mut self) {
946 if Rc::strong_count(&self.inner) != 1 {
947 return;
948 }
949 unregister_runtime_handle(self.inner.runtime_id);
950 LAST_RUNTIME.with(|slot| {
951 let should_clear = slot
952 .borrow()
953 .as_ref()
954 .is_some_and(|handle| handle.id() == self.inner.runtime_id);
955 if should_clear {
956 *slot.borrow_mut() = None;
957 }
958 });
959 }
960}
961
962#[derive(Default)]
963pub struct DefaultScheduler;
964
965impl RuntimeScheduler for DefaultScheduler {
966 fn schedule_frame(&self) {}
967}
968
969#[cfg(test)]
970#[derive(Default)]
971pub struct TestScheduler;
972
973#[cfg(test)]
974impl RuntimeScheduler for TestScheduler {
975 fn schedule_frame(&self) {}
976}
977
978#[cfg(test)]
979pub struct TestRuntime {
980 runtime: Runtime,
981}
982
983#[cfg(test)]
984impl Default for TestRuntime {
985 fn default() -> Self {
986 Self::new()
987 }
988}
989
990#[cfg(test)]
991impl TestRuntime {
992 pub fn new() -> Self {
993 Self {
994 runtime: Runtime::new(Arc::new(TestScheduler)),
995 }
996 }
997
998 pub fn handle(&self) -> RuntimeHandle {
999 self.runtime.handle()
1000 }
1001}
1002
1003#[derive(Clone)]
1004pub struct RuntimeHandle {
1005 inner: Weak<RuntimeInner>,
1006 dispatcher: UiDispatcher,
1007 ui_thread_id: ThreadId,
1008 id: RuntimeId,
1009}
1010
1011pub(crate) struct ScopeRuntime(Weak<RuntimeInner>);
1015
1016impl ScopeRuntime {
1017 pub(crate) fn mark_scope_recomposed(&self) {
1018 if let Some(inner) = self.0.upgrade() {
1019 inner.mark_scope_recomposed();
1020 }
1021 }
1022
1023 pub(crate) fn register_invalid_scope(&self, scope: Weak<RecomposeScopeInner>) {
1024 if let Some(inner) = self.0.upgrade() {
1025 inner.register_invalid_scope(scope);
1026 }
1027 }
1028
1029 pub(crate) fn unregister_state_scope(&self, id: StateId, scope_id: ScopeId) {
1030 if let Some(inner) = self.0.upgrade() {
1031 inner.state_arena.unregister_scope(id, scope_id);
1032 }
1033 }
1034
1035 pub(crate) fn increment_live_recompose_scope_count(&self) {
1036 if let Some(inner) = self.0.upgrade() {
1037 inner.increment_live_recompose_scope_count();
1038 }
1039 }
1040
1041 pub(crate) fn decrement_live_recompose_scope_count(&self) {
1042 if let Some(inner) = self.0.upgrade() {
1043 inner.decrement_live_recompose_scope_count();
1044 }
1045 }
1046}
1047
1048pub struct TaskHandle {
1049 id: u64,
1050 runtime: RuntimeHandle,
1051}
1052
1053struct DeferredStateRelease {
1054 runtime: RuntimeHandle,
1055 id: StateId,
1056}
1057
1058pub(crate) struct StateHandleLease {
1059 id: StateId,
1060 runtime: RuntimeHandle,
1061}
1062
1063impl StateHandleLease {
1064 pub(crate) fn id(&self) -> StateId {
1065 self.id
1066 }
1067
1068 pub(crate) fn runtime(&self) -> RuntimeHandle {
1069 self.runtime.clone()
1070 }
1071}
1072
1073impl Drop for StateHandleLease {
1074 fn drop(&mut self) {
1075 defer_state_release(self.runtime.clone(), self.id);
1076 }
1077}
1078
1079thread_local! {
1080 static STATE_OWNERS: RefCell<Vec<Vec<Rc<StateHandleLease>>>> =
1081 const { RefCell::new(Vec::new()) };
1082}
1083
1084struct StateOwnerFrame;
1085
1086impl Drop for StateOwnerFrame {
1087 fn drop(&mut self) {
1088 STATE_OWNERS.with(|owners| owners.borrow_mut().pop());
1089 }
1090}
1091
1092pub(crate) fn collecting_states<T>(build: impl FnOnce() -> T) -> (T, Vec<Rc<StateHandleLease>>) {
1110 STATE_OWNERS.with(|owners| owners.borrow_mut().push(Vec::new()));
1111 let frame = StateOwnerFrame;
1112 let value = build();
1113 let states = STATE_OWNERS.with(|owners| {
1114 owners
1115 .borrow_mut()
1116 .last_mut()
1117 .map(std::mem::take)
1118 .unwrap_or_default()
1119 });
1120 drop(frame);
1121 (value, states)
1122}
1123
1124fn hand_to_current_owner(lease: &Rc<StateHandleLease>) -> bool {
1125 STATE_OWNERS.with(|owners| match owners.borrow_mut().last_mut() {
1126 Some(owner) => {
1127 owner.push(Rc::clone(lease));
1128 true
1129 }
1130 None => false,
1131 })
1132}
1133
1134impl RuntimeHandle {
1135 pub fn id(&self) -> RuntimeId {
1136 self.id
1137 }
1138
1139 pub(crate) fn alloc_state<T: Clone + 'static>(&self, value: T) -> Rc<StateHandleLease> {
1140 let id = self.with_state_arena(|arena| arena.alloc(value, self.clone()));
1141 let lease = Rc::new(StateHandleLease {
1142 id,
1143 runtime: self.clone(),
1144 });
1145 self.with_state_arena(|arena| arena.register_lease(id, &lease));
1146 lease
1147 }
1148
1149 pub(crate) fn alloc_state_with_policy<T: Clone + 'static>(
1150 &self,
1151 value: T,
1152 policy: Rc<dyn MutationPolicy<T>>,
1153 ) -> Rc<StateHandleLease> {
1154 let id =
1155 self.with_state_arena(|arena| arena.alloc_with_policy(value, self.clone(), policy));
1156 let lease = Rc::new(StateHandleLease {
1157 id,
1158 runtime: self.clone(),
1159 });
1160 self.with_state_arena(|arena| arena.register_lease(id, &lease));
1161 lease
1162 }
1163
1164 pub(crate) fn alloc_persistent_state<T: Clone + 'static>(
1165 &self,
1166 value: T,
1167 ) -> crate::MutableState<T> {
1168 self.hand_out(self.alloc_state(value))
1169 }
1170
1171 pub(crate) fn alloc_persistent_state_with_policy<T: Clone + 'static>(
1172 &self,
1173 value: T,
1174 policy: Rc<dyn MutationPolicy<T>>,
1175 ) -> crate::MutableState<T> {
1176 self.hand_out(self.alloc_state_with_policy(value, policy))
1177 }
1178
1179 fn hand_out<T: Clone + 'static>(&self, lease: Rc<StateHandleLease>) -> crate::MutableState<T> {
1180 if !hand_to_current_owner(&lease)
1181 && let Some(inner) = self.inner.upgrade()
1182 {
1183 inner
1184 .external_state_owners
1185 .borrow_mut()
1186 .insert(lease.id(), Rc::clone(&lease));
1187 }
1188 crate::MutableState::from_lease(&lease)
1189 }
1190
1191 pub(crate) fn retain_state_lease(&self, id: StateId) -> Option<Rc<StateHandleLease>> {
1192 self.with_state_arena(|arena| arena.retain_lease(id))
1193 }
1194
1195 pub(crate) fn with_state_arena<R>(&self, f: impl FnOnce(&StateArena) -> R) -> R {
1196 self.try_with_state_arena(f)
1197 .unwrap_or_else(|| panic!("runtime dropped"))
1198 }
1199
1200 pub(crate) fn try_with_state_arena<R>(&self, f: impl FnOnce(&StateArena) -> R) -> Option<R> {
1201 self.inner.upgrade().map(|inner| f(&inner.state_arena))
1202 }
1203
1204 fn release_state_immediate(&self, id: StateId) {
1205 if let Some(inner) = self.inner.upgrade() {
1206 inner.state_arena.release(id);
1207 }
1208 }
1209
1210 pub fn state_arena_stats(&self) -> (usize, usize) {
1211 self.try_with_state_arena(StateArena::stats)
1212 .unwrap_or_default()
1213 }
1214
1215 pub fn state_arena_debug_stats(&self) -> StateArenaDebugStats {
1216 self.try_with_state_arena(StateArena::debug_stats)
1217 .unwrap_or_default()
1218 }
1219
1220 pub fn debug_stats(&self) -> RuntimeDebugStats {
1221 self.inner
1222 .upgrade()
1223 .map(|inner| inner.debug_stats())
1224 .unwrap_or_default()
1225 }
1226
1227 pub fn live_ui_task_labels(&self) -> Vec<(u64, String)> {
1228 self.inner
1229 .upgrade()
1230 .map(|inner| {
1231 inner
1232 .tasks
1233 .borrow()
1234 .iter()
1235 .map(|(id, entry)| (*id, entry.label.clone()))
1236 .collect()
1237 })
1238 .unwrap_or_default()
1239 }
1240
1241 pub fn schedule(&self) {
1242 if let Some(inner) = self.inner.upgrade() {
1243 inner.schedule();
1244 }
1245 }
1246
1247 pub(crate) fn enqueue_node_update(&self, command: Command) {
1248 if let Some(inner) = self.inner.upgrade() {
1249 inner.enqueue_update(command);
1250 }
1251 }
1252
1253 pub fn enqueue_ui_task(&self, task: Box<dyn FnOnce() + 'static>) {
1260 if let Some(inner) = self.inner.upgrade() {
1261 inner.enqueue_ui_task(task);
1262 } else {
1263 task();
1264 }
1265 }
1266
1267 pub fn spawn_ui<F>(&self, fut: F) -> Option<TaskHandle>
1268 where
1269 F: Future<Output = ()> + 'static,
1270 {
1271 self.inner.upgrade().map(|inner| {
1272 let id = inner.spawn_ui_task(Box::pin(fut));
1273 TaskHandle {
1274 id,
1275 runtime: self.clone(),
1276 }
1277 })
1278 }
1279
1280 pub fn cancel_task(&self, id: u64) {
1281 if let Some(inner) = self.inner.upgrade() {
1282 inner.cancel_task(id);
1283 }
1284 }
1285
1286 pub fn has_task(&self, id: u64) -> bool {
1288 self.inner.upgrade().is_some_and(|inner| inner.has_task(id))
1289 }
1290
1291 pub fn post_ui(&self, task: impl FnOnce() + Send + 'static) {
1296 self.dispatcher.post(task);
1297 }
1298
1299 pub fn register_ui_cont<T: 'static>(&self, f: impl FnOnce(T) + 'static) -> Option<u64> {
1300 self.inner.upgrade().map(|inner| inner.register_ui_cont(f))
1301 }
1302
1303 pub fn cancel_ui_cont(&self, id: u64) {
1304 if let Some(inner) = self.inner.upgrade() {
1305 inner.cancel_ui_cont(id);
1306 }
1307 }
1308
1309 pub fn drain_ui(&self) {
1310 if let Some(inner) = self.inner.upgrade() {
1311 inner.drain_ui();
1312 }
1313 }
1314
1315 pub fn has_pending_ui(&self) -> bool {
1316 self.inner.upgrade().map_or_else(
1317 || self.dispatcher.has_pending(),
1318 |inner| inner.has_pending_ui(),
1319 )
1320 }
1321
1322 pub fn register_frame_callback(
1323 &self,
1324 callback: impl FnOnce(u64) + 'static,
1325 ) -> Option<FrameCallbackId> {
1326 self.inner.upgrade().map(|inner| {
1327 inner.register_frame_callback(FrameCallbackKind::Transient, Box::new(callback))
1328 })
1329 }
1330
1331 pub fn register_perpetual_frame_callback(
1332 &self,
1333 callback: impl FnOnce(u64) + 'static,
1334 ) -> Option<FrameCallbackId> {
1335 self.inner.upgrade().map(|inner| {
1336 inner.register_frame_callback(FrameCallbackKind::Perpetual, Box::new(callback))
1337 })
1338 }
1339
1340 pub fn cancel_frame_callback(&self, id: FrameCallbackId) {
1341 if let Some(inner) = self.inner.upgrade() {
1342 inner.cancel_frame_callback(id);
1343 }
1344 }
1345
1346 pub fn drain_frame_callbacks(&self, frame_time_nanos: u64) {
1350 if let Some(inner) = self.inner.upgrade() {
1351 inner.drain_frame_callbacks(frame_time_nanos);
1352 }
1353 }
1354
1355 pub fn last_frame_time_nanos(&self) -> Option<u64> {
1358 self.inner
1359 .upgrade()
1360 .and_then(|inner| inner.last_frame_time_nanos.get())
1361 }
1362
1363 #[cfg(any(feature = "internal", test))]
1364 pub fn frame_clock(&self) -> FrameClock {
1365 FrameClock::new(self.clone())
1366 }
1367
1368 pub fn set_needs_frame(&self, value: bool) {
1369 if let Some(inner) = self.inner.upgrade() {
1370 *inner.needs_frame.borrow_mut() = value;
1371 }
1372 }
1373
1374 pub(crate) fn take_updates(&self) -> Vec<Command> {
1375 self.inner
1376 .upgrade()
1377 .map(|inner| inner.take_updates())
1378 .unwrap_or_default()
1379 }
1380
1381 pub fn has_updates(&self) -> bool {
1382 self.inner
1383 .upgrade()
1384 .is_some_and(|inner| inner.has_updates())
1385 }
1386
1387 pub(crate) fn scope_runtime(&self) -> ScopeRuntime {
1389 ScopeRuntime(Weak::clone(&self.inner))
1390 }
1391
1392 pub(crate) fn requeue_invalid_scope(&self, scope: &RecomposeScope) {
1393 if let Some(inner) = self.inner.upgrade() {
1394 inner.requeue_invalid_scope(scope);
1395 }
1396 }
1397
1398 pub(crate) fn take_invalidated_scopes(&self) -> Option<Vec<RecomposeScope>> {
1399 self.inner
1400 .upgrade()
1401 .and_then(|inner| inner.take_invalidated_scopes())
1402 }
1403
1404 pub fn forget_movable(&self, id: Key) {
1408 if let Some(inner) = self.inner.upgrade() {
1409 inner.forgotten_movables.borrow_mut().push(id);
1410 inner.schedule();
1411 }
1412 }
1413
1414 pub(crate) fn take_forgotten_movables(&self) -> Vec<Key> {
1415 self.inner
1416 .upgrade()
1417 .map(|inner| std::mem::take(&mut *inner.forgotten_movables.borrow_mut()))
1418 .unwrap_or_default()
1419 }
1420
1421 pub(crate) fn next_movable_content_id(&self) -> Key {
1426 let Some(inner) = self.inner.upgrade() else {
1427 log::error!("movable content asked a runtime that is gone for an identity");
1428 return 0;
1429 };
1430 let id = inner.next_movable_content_id.get();
1431 inner.next_movable_content_id.set(id.wrapping_add(1).max(1));
1432 id
1433 }
1434
1435 pub fn has_invalid_scopes(&self) -> bool {
1436 self.inner
1437 .upgrade()
1438 .is_some_and(|inner| inner.has_invalid_scopes())
1439 }
1440
1441 fn live_recompose_scope_count(&self) -> usize {
1442 self.inner
1443 .upgrade()
1444 .map(|inner| inner.live_recompose_scope_count())
1445 .unwrap_or_default()
1446 }
1447
1448 #[doc(hidden)]
1449 pub fn debug_invalid_scope_ids(&self) -> Vec<usize> {
1450 self.inner
1451 .upgrade()
1452 .map(|inner| inner.queued_invalid_scope_ids())
1453 .unwrap_or_default()
1454 }
1455
1456 pub fn has_frame_callbacks(&self) -> bool {
1457 self.inner
1458 .upgrade()
1459 .is_some_and(|inner| inner.has_frame_callbacks())
1460 }
1461
1462 pub fn has_transient_frame_callbacks(&self) -> bool {
1463 self.inner
1464 .upgrade()
1465 .is_some_and(|inner| inner.has_transient_frame_callbacks())
1466 }
1467
1468 pub fn assert_ui_thread(&self) {
1469 debug_assert_eq!(
1470 std::thread::current().id(),
1471 self.ui_thread_id,
1472 "state mutated off the runtime's UI thread"
1473 );
1474 }
1475
1476 pub fn dispatcher(&self) -> UiDispatcher {
1477 self.dispatcher.clone()
1478 }
1479
1480 #[doc(hidden)]
1481 pub fn with_deferred_state_releases<R>(&self, f: impl FnOnce() -> R) -> R {
1482 let _scope = enter_state_teardown_scope();
1483 f()
1484 }
1485}
1486
1487impl TaskHandle {
1488 pub fn cancel(&self) {
1489 self.runtime.cancel_task(self.id);
1490 }
1491
1492 pub fn is_finished(&self) -> bool {
1494 !self.runtime.has_task(self.id)
1495 }
1496}
1497
1498pub(crate) struct FrameCallbackEntry {
1499 id: FrameCallbackId,
1500 kind: FrameCallbackKind,
1501 callback: Option<Box<dyn FnOnce(u64) + 'static>>,
1502}
1503
1504#[cfg(not(target_arch = "wasm32"))]
1505struct RuntimeTaskWaker {
1506 scheduler: SchedulerRef,
1507 runnable: Arc<AtomicBool>,
1508}
1509
1510#[cfg(target_arch = "wasm32")]
1511struct RuntimeTaskWaker {
1512 runtime_id: RuntimeId,
1513 runnable: Arc<AtomicBool>,
1514}
1515
1516impl RuntimeTaskWaker {
1517 #[cfg(not(target_arch = "wasm32"))]
1518 fn new(inner: &RuntimeInner, runnable: Arc<AtomicBool>) -> Self {
1519 let scheduler = inner.scheduler.clone();
1520 Self {
1521 scheduler,
1522 runnable,
1523 }
1524 }
1525
1526 #[cfg(target_arch = "wasm32")]
1527 fn new(inner: &RuntimeInner, runnable: Arc<AtomicBool>) -> Self {
1528 let runtime_id = inner.runtime_id;
1529 Self {
1530 runtime_id,
1531 runnable,
1532 }
1533 }
1534
1535 fn into_waker(self) -> Waker {
1536 futures_task::waker(Arc::new(self))
1537 }
1538}
1539
1540impl futures_task::ArcWake for RuntimeTaskWaker {
1541 #[cfg(not(target_arch = "wasm32"))]
1542 fn wake_by_ref(arc_self: &Arc<Self>) {
1543 arc_self.runnable.store(true, Ordering::Release);
1544 arc_self.scheduler.schedule_frame();
1545 }
1546
1547 #[cfg(target_arch = "wasm32")]
1548 fn wake_by_ref(arc_self: &Arc<Self>) {
1549 arc_self.runnable.store(true, Ordering::Release);
1550 REGISTERED_RUNTIMES.with(|registry| {
1551 if let Some(handle) = registry.borrow().get(&arc_self.runtime_id).cloned() {
1552 handle.schedule();
1553 }
1554 });
1555 }
1556}
1557
1558thread_local! {
1559 static NEXT_RUNTIME_ID: Cell<u32> = const { Cell::new(1) };
1560 static ACTIVE_RUNTIMES: RefCell<Vec<RuntimeHandle>> = const { RefCell::new(Vec::new()) };
1561 static LAST_RUNTIME: RefCell<Option<RuntimeHandle>> = const { RefCell::new(None) };
1562 static REGISTERED_RUNTIMES: RefCell<HashMap<RuntimeId, RuntimeHandle>> = RefCell::new(HashMap::default());
1563 static STATE_TEARDOWN_DEPTH: Cell<usize> = const { Cell::new(0) };
1564 static DEFERRED_STATE_RELEASES: RefCell<Vec<DeferredStateRelease>> = const { RefCell::new(Vec::new()) };
1565}
1566
1567#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1568pub struct RuntimeThreadLocalDebugStats {
1569 pub active_runtimes_len: usize,
1570 pub active_runtimes_cap: usize,
1571 pub registered_runtimes_len: usize,
1572 pub registered_runtimes_cap: usize,
1573 pub deferred_state_releases_len: usize,
1574 pub deferred_state_releases_cap: usize,
1575}
1576
1577pub fn current_runtime_handle() -> Option<RuntimeHandle> {
1582 if let Some(handle) = ACTIVE_RUNTIMES.with(|stack| stack.borrow().last().cloned()) {
1583 return Some(handle);
1584 }
1585 LAST_RUNTIME.with(|slot| slot.borrow().clone())
1586}
1587
1588pub(crate) fn runtime_handle_by_id(id: RuntimeId) -> Option<RuntimeHandle> {
1589 REGISTERED_RUNTIMES.with(|registry| registry.borrow().get(&id).cloned())
1590}
1591
1592pub(crate) fn with_state_arena_by_id<R>(
1593 id: RuntimeId,
1594 f: impl FnOnce(&StateArena) -> R,
1595) -> Option<R> {
1596 let inner = REGISTERED_RUNTIMES.with(|registry| {
1597 registry
1598 .borrow()
1599 .get(&id)
1600 .and_then(|handle| handle.inner.upgrade())
1601 })?;
1602 Some(f(&inner.state_arena))
1603}
1604
1605pub(crate) fn live_recompose_scope_count() -> usize {
1606 REGISTERED_RUNTIMES.with(|registry| {
1607 registry
1608 .borrow()
1609 .values()
1610 .map(RuntimeHandle::live_recompose_scope_count)
1611 .sum()
1612 })
1613}
1614
1615pub fn debug_runtime_thread_local_stats() -> RuntimeThreadLocalDebugStats {
1616 let (active_runtimes_len, active_runtimes_cap) = ACTIVE_RUNTIMES.with(|stack| {
1617 let stack = stack.borrow();
1618 (stack.len(), stack.capacity())
1619 });
1620 let (registered_runtimes_len, registered_runtimes_cap) = REGISTERED_RUNTIMES.with(|registry| {
1621 let registry = registry.borrow();
1622 (registry.len(), registry.capacity())
1623 });
1624 let (deferred_state_releases_len, deferred_state_releases_cap) =
1625 DEFERRED_STATE_RELEASES.with(|releases| {
1626 let releases = releases.borrow();
1627 (releases.len(), releases.capacity())
1628 });
1629
1630 RuntimeThreadLocalDebugStats {
1631 active_runtimes_len,
1632 active_runtimes_cap,
1633 registered_runtimes_len,
1634 registered_runtimes_cap,
1635 deferred_state_releases_len,
1636 deferred_state_releases_cap,
1637 }
1638}
1639
1640fn register_runtime_handle(handle: &RuntimeHandle) {
1641 REGISTERED_RUNTIMES.with(|registry| {
1642 registry.borrow_mut().insert(handle.id(), handle.clone());
1643 });
1644}
1645
1646fn unregister_runtime_handle(id: RuntimeId) {
1647 REGISTERED_RUNTIMES.with(|registry| {
1648 registry.borrow_mut().remove(&id);
1649 });
1650}
1651
1652fn defer_state_release(runtime: RuntimeHandle, id: StateId) {
1653 let teardown_active = STATE_TEARDOWN_DEPTH.with(|depth| depth.get() > 0);
1654 if teardown_active {
1655 DEFERRED_STATE_RELEASES.with(|releases| {
1656 releases
1657 .borrow_mut()
1658 .push(DeferredStateRelease { runtime, id });
1659 });
1660 } else {
1661 runtime.release_state_immediate(id);
1662 }
1663}
1664
1665fn flush_deferred_state_releases() {
1666 DEFERRED_STATE_RELEASES.with(|releases| {
1667 let mut releases = releases.borrow_mut();
1668 while let Some(deferred) = releases.pop() {
1669 deferred.runtime.release_state_immediate(deferred.id);
1670 }
1671 });
1672}
1673
1674pub(crate) struct StateTeardownScope;
1675
1676pub(crate) fn enter_state_teardown_scope() -> StateTeardownScope {
1677 STATE_TEARDOWN_DEPTH.with(|depth| depth.set(depth.get() + 1));
1678 StateTeardownScope
1679}
1680
1681impl Drop for StateTeardownScope {
1682 fn drop(&mut self) {
1683 STATE_TEARDOWN_DEPTH.with(|depth| {
1684 let next = depth.get().saturating_sub(1);
1685 depth.set(next);
1686 if next == 0 {
1687 flush_deferred_state_releases();
1688 }
1689 });
1690 }
1691}
1692
1693pub(crate) fn push_active_runtime(handle: &RuntimeHandle) {
1694 register_runtime_handle(handle);
1695 ACTIVE_RUNTIMES.with(|stack| stack.borrow_mut().push(handle.clone()));
1696 LAST_RUNTIME.with(|slot| *slot.borrow_mut() = Some(handle.clone()));
1697}
1698
1699pub(crate) fn pop_active_runtime() {
1700 ACTIVE_RUNTIMES.with(|stack| {
1701 stack.borrow_mut().pop();
1702 });
1703}
1704
1705pub fn schedule_frame() {
1707 if let Some(handle) = current_runtime_handle() {
1708 handle.schedule();
1709 return;
1710 }
1711 log::debug!(
1712 target: "cranpose::runtime",
1713 "ignoring frame request without an active runtime",
1714 );
1715}
1716
1717pub fn schedule_node_update(
1719 update: impl FnOnce(&mut dyn Applier) -> Result<(), NodeError> + 'static,
1720) {
1721 if let Some(handle) = current_runtime_handle() {
1722 handle.enqueue_node_update(Command::callback(update));
1723 } else {
1724 drop(update);
1725 log::debug!(
1726 target: "cranpose::runtime",
1727 "ignoring node update request without an active runtime",
1728 );
1729 }
1730}
1731
1732#[cfg(test)]
1733#[path = "tests/runtime_tests.rs"]
1734mod tests;