Skip to main content

cranpose_core/
runtime.rs

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