Skip to main content

cranpose_core/
snapshot_state_observer.rs

1use std::{
2    any::{Any, TypeId},
3    cell::{Cell, RefCell},
4    hash::{Hash, Hasher},
5    rc::{Rc, Weak},
6    sync::Arc,
7};
8
9use smallvec::SmallVec;
10
11use crate::{
12    collections::map::{HashMap, HashSet},
13    hash::default as default_hash,
14    snapshot_v2::{
15        ReadObserver, StateObjectId, TransparentObserverMutableSnapshot, register_apply_observer,
16    },
17    state::StateObject,
18};
19
20type Executor = dyn Fn(Box<dyn FnOnce() + 'static>) + 'static;
21
22trait ScopeChangedCallback: Fn(&dyn Any) + Any {}
23
24impl<F: Fn(&dyn Any) + Any> ScopeChangedCallback for F {}
25
26/// Observer that records state object reads performed inside a given scope and
27/// notifies the caller when any of the observed objects change.
28///
29/// This is a pragmatic Rust translation of Jetpack Compose's
30/// `SnapshotStateObserver`. The implementation focuses on the core behaviour
31/// needed by the Cranpose runtime:
32/// - Tracking state object reads per logical scope.
33/// - Reacting to snapshot apply notifications.
34/// - Scheduling invalidation callbacks via the supplied executor.
35///
36/// Advanced features from the Kotlin version (derived state tracking, change
37/// coalescing, queue minimisation) are deferred
38#[derive(Clone)]
39pub struct SnapshotStateObserver {
40    inner: Rc<SnapshotStateObserverInner>,
41}
42
43#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
44pub struct SnapshotStateObserverDebugStats {
45    pub scopes_len: usize,
46    pub scopes_cap: usize,
47    pub stateless_scope_count: usize,
48    pub observed_state_count: usize,
49    pub observed_state_capacity: usize,
50}
51
52impl SnapshotStateObserver {
53    /// Create a new observer that schedules callbacks using `on_changed_executor`.
54    pub fn new(on_changed_executor: impl Fn(Box<dyn FnOnce() + 'static>) + 'static) -> Self {
55        Self {
56            inner: Rc::new(SnapshotStateObserverInner::new(on_changed_executor)),
57        }
58    }
59
60    /// Observe state object reads performed while executing `block`.
61    ///
62    /// Subsequent calls to `observe_reads` replace any previously recorded
63    /// observations for the provided `scope`. When one of the observed objects
64    /// mutates, `on_value_changed_for_scope` will be invoked on the executor.
65    pub fn observe_reads<T, R>(
66        &self,
67        scope: T,
68        on_value_changed_for_scope: impl Fn(&T) + 'static,
69        block: impl FnOnce() -> R,
70    ) -> R
71    where
72        T: Any + Clone + Eq + Hash + 'static,
73    {
74        self.inner
75            .observe_reads(scope, on_value_changed_for_scope, block)
76    }
77
78    /// Temporarily pause read observation while executing `block`.
79    pub fn with_no_observations<R>(&self, block: impl FnOnce() -> R) -> R {
80        self.inner.with_no_observations(block)
81    }
82
83    /// Remove any recorded reads for `scope`.
84    pub fn clear<T>(&self, scope: &T)
85    where
86        T: Any + Eq + Hash + 'static,
87    {
88        self.inner.clear(scope);
89    }
90
91    /// Remove recorded reads for scopes that satisfy `predicate`.
92    pub fn clear_if(&self, predicate: impl Fn(&dyn Any) -> bool) {
93        self.inner.clear_if(predicate);
94    }
95
96    /// Remove all recorded observations.
97    pub fn clear_all(&self) {
98        self.inner.clear_all();
99    }
100
101    /// Begin listening for snapshot apply notifications.
102    pub fn start(&self) {
103        let weak = Rc::downgrade(&self.inner);
104        self.inner.start(weak);
105    }
106
107    /// Stop listening for snapshot apply notifications.
108    pub fn stop(&self) {
109        self.inner.stop();
110    }
111
112    pub fn debug_stats(&self) -> SnapshotStateObserverDebugStats {
113        self.inner.debug_stats()
114    }
115
116    #[cfg(test)]
117    pub fn notify_changes(&self, modified: &[Arc<dyn StateObject>]) {
118        self.inner.handle_apply(modified);
119    }
120}
121
122struct SnapshotStateObserverInner {
123    executor: Rc<Executor>,
124    owned_scopes: RefCell<HashMap<OwnedScopeIndexKey, OwnedScopeBucket>>,
125    indexed_scopes: RefCell<HashMap<usize, Rc<RefCell<ScopeEntry>>>>,
126    observed_to_scopes: RefCell<HashMap<StateObjectId, HashSet<usize>>>,
127    pause_count: Rc<Cell<usize>>,
128    active_read_targets: Rc<RefCell<ReadObservationStack>>,
129    read_dispatcher: ReadObserver,
130    read_snapshot: RefCell<Option<Arc<TransparentObserverMutableSnapshot>>>,
131    apply_handle: RefCell<Option<crate::snapshot_v2::ObserverHandle>>,
132    next_entry_id: Cell<usize>,
133    /// One `Rc` per type of callback that captures nothing: see
134    /// [`SnapshotStateObserverInner::capture_free_callback`].
135    capture_free_callbacks: RefCell<CaptureFreeCallbacks>,
136}
137
138#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
139struct OwnedScopeIndexKey {
140    type_id: TypeId,
141    value_hash: u64,
142}
143
144type OwnedScopeBucket = SmallVec<[Rc<RefCell<ScopeEntry>>; 1]>;
145
146/// The shared `Rc` of each type of callback that captures nothing.
147type CaptureFreeCallbacks = SmallVec<[(TypeId, Rc<dyn ScopeChangedCallback>); 2]>;
148
149fn owned_scope_index_key<T>(scope: &T) -> OwnedScopeIndexKey
150where
151    T: Any + Hash + 'static,
152{
153    let mut hasher = default_hash::new();
154    scope.hash(&mut hasher);
155    OwnedScopeIndexKey {
156        type_id: TypeId::of::<T>(),
157        value_hash: hasher.finish(),
158    }
159}
160
161impl SnapshotStateObserverInner {
162    const MIN_RETAINED_SCOPE_CAPACITY: usize = 256;
163
164    fn new(on_changed_executor: impl Fn(Box<dyn FnOnce() + 'static>) + 'static) -> Self {
165        let pause_count = Rc::new(Cell::new(0));
166        let active_read_targets = Rc::new(RefCell::new(ReadObservationStack::default()));
167        let dispatcher_pause_count = Rc::clone(&pause_count);
168        let dispatcher_targets = Rc::clone(&active_read_targets);
169        let read_dispatcher: ReadObserver = Arc::new(move |state| {
170            if dispatcher_pause_count.get() > 0 {
171                return;
172            }
173            let observed = dispatcher_targets.borrow().last().cloned();
174            if let Some(observed) = observed {
175                observed.borrow_mut().insert(state);
176            }
177        });
178
179        Self {
180            executor: Rc::new(on_changed_executor),
181            owned_scopes: RefCell::new(HashMap::default()),
182            indexed_scopes: RefCell::new(HashMap::default()),
183            observed_to_scopes: RefCell::new(HashMap::default()),
184            pause_count,
185            active_read_targets,
186            read_dispatcher,
187            read_snapshot: RefCell::new(None),
188            apply_handle: RefCell::new(None),
189            next_entry_id: Cell::new(0),
190            capture_free_callbacks: RefCell::new(SmallVec::new()),
191        }
192    }
193
194    fn observe_reads<T, R>(
195        &self,
196        scope: T,
197        on_value_changed_for_scope: impl Fn(&T) + 'static,
198        block: impl FnOnce() -> R,
199    ) -> R
200    where
201        T: Any + Clone + Eq + Hash + 'static,
202    {
203        let existing_entry = self.find_scope_entry(&scope);
204        let on_changed = std::cell::LazyCell::new(|| {
205            let callback = move |scope_any: &dyn Any| {
206                if let Some(typed) = scope_any.downcast_ref::<T>() {
207                    on_value_changed_for_scope(typed);
208                }
209            };
210            if std::mem::size_of_val(&callback) == 0 {
211                return self.capture_free_callback(callback);
212            }
213            match existing_entry.as_ref() {
214                Some(entry) => entry.borrow_mut().callback_reusing(callback),
215                None => Rc::new(callback),
216            }
217        });
218
219        if let Some(entry) = existing_entry.as_ref() {
220            entry.borrow_mut().update_scope(&scope);
221            let callback = on_changed.clone();
222            entry.borrow_mut().on_changed = callback;
223        }
224
225        let observed = self.active_read_targets.borrow_mut().push();
226        struct ActiveObservationGuard<'a> {
227            stack: &'a RefCell<ReadObservationStack>,
228        }
229        impl Drop for ActiveObservationGuard<'_> {
230            fn drop(&mut self) {
231                let target = self.stack.borrow_mut().pop();
232                target.borrow_mut().clear();
233            }
234        }
235        let _guard = ActiveObservationGuard {
236            stack: &self.active_read_targets,
237        };
238
239        let result = self.run_with_read_observer(block);
240
241        if observed.borrow().is_empty() {
242            if existing_entry.is_some() {
243                self.clear(&scope);
244            }
245            return result;
246        }
247
248        let callback = Rc::clone(&on_changed);
249        drop(on_changed);
250        let entry = existing_entry
251            .unwrap_or_else(|| self.insert_scope_entry(scope.clone(), Rc::clone(&callback)));
252        entry.borrow_mut().update(&scope, callback);
253        self.replace_observed_ids(&entry, &mut observed.borrow_mut());
254
255        result
256    }
257
258    fn with_no_observations<R>(&self, block: impl FnOnce() -> R) -> R {
259        self.pause_count.set(self.pause_count.get() + 1);
260        let result = block();
261        self.pause_count
262            .set(self.pause_count.get().saturating_sub(1));
263        result
264    }
265
266    fn clear<T>(&self, scope: &T)
267    where
268        T: Any + Eq + Hash + 'static,
269    {
270        let removed = self.remove_scope_entry(scope);
271        if let Some(entry) = removed {
272            self.unregister_entry(&entry);
273        }
274    }
275
276    fn clear_if(&self, predicate: impl Fn(&dyn Any) -> bool) {
277        let removed = self.partition_scopes(|entry| predicate(entry.scope.as_ref()));
278        for entry in removed {
279            self.unregister_entry(&entry);
280        }
281    }
282
283    fn clear_all(&self) {
284        let entries = std::mem::take(&mut *self.indexed_scopes.borrow_mut());
285        let owned = std::mem::take(&mut *self.owned_scopes.borrow_mut());
286        self.observed_to_scopes.borrow_mut().clear();
287        drop((entries, owned));
288    }
289
290    fn start(&self, weak_self: Weak<SnapshotStateObserverInner>) {
291        if self.apply_handle.borrow().is_some() {
292            return;
293        }
294
295        let handle = register_apply_observer(Rc::new(move |modified, _snapshot_id| {
296            if let Some(inner) = weak_self.upgrade() {
297                inner.handle_apply(modified);
298            }
299        }));
300        self.apply_handle.replace(Some(handle));
301    }
302
303    fn stop(&self) {
304        if let Some(handle) = self.apply_handle.borrow_mut().take() {
305            drop(handle);
306        }
307    }
308
309    /// The `Rc` of a callback that captures nothing. Every closure of such a
310    /// type does the same thing, so one `Rc` serves all its scopes instead
311    /// of one each.
312    fn capture_free_callback<F: Fn(&dyn Any) + 'static>(
313        &self,
314        callback: F,
315    ) -> Rc<dyn ScopeChangedCallback> {
316        let type_id = TypeId::of::<F>();
317        let mut shared = self.capture_free_callbacks.borrow_mut();
318        if let Some((_, callback)) = shared.iter().find(|(id, _)| *id == type_id) {
319            return Rc::clone(callback);
320        }
321        let callback: Rc<dyn ScopeChangedCallback> = Rc::new(callback);
322        shared.push((type_id, Rc::clone(&callback)));
323        callback
324    }
325
326    fn insert_scope_entry(
327        &self,
328        scope: impl Any + Clone + Eq + Hash + 'static,
329        on_changed: Rc<dyn ScopeChangedCallback>,
330    ) -> Rc<RefCell<ScopeEntry>> {
331        let entry_id = self.next_entry_id.get();
332        self.next_entry_id.set(entry_id.wrapping_add(1));
333        let scope_key = owned_scope_index_key(&scope);
334        let entry = Rc::new(RefCell::new(ScopeEntry::new(entry_id, scope, on_changed)));
335        self.indexed_scopes
336            .borrow_mut()
337            .insert(entry_id, Rc::clone(&entry));
338        self.owned_scopes
339            .borrow_mut()
340            .entry(scope_key)
341            .or_default()
342            .push(Rc::clone(&entry));
343        entry
344    }
345
346    fn find_scope_entry<T>(&self, scope: &T) -> Option<Rc<RefCell<ScopeEntry>>>
347    where
348        T: Any + Eq + Hash + 'static,
349    {
350        let key = owned_scope_index_key(scope);
351        self.owned_scopes.borrow().get(&key).and_then(|bucket| {
352            bucket
353                .iter()
354                .find(|entry| entry.borrow().matches_scope(scope))
355                .cloned()
356        })
357    }
358
359    fn remove_scope_entry<T>(&self, scope: &T) -> Option<Rc<RefCell<ScopeEntry>>>
360    where
361        T: Any + Eq + Hash + 'static,
362    {
363        let key = owned_scope_index_key(scope);
364        let mut owned_scopes = self.owned_scopes.borrow_mut();
365        let mut removed = None;
366        let mut remove_bucket = false;
367        if let Some(bucket) = owned_scopes.get_mut(&key)
368            && let Some(index) = bucket
369                .iter()
370                .position(|entry| entry.borrow().matches_scope(scope))
371        {
372            removed = Some(bucket.remove(index));
373            remove_bucket = bucket.is_empty();
374        }
375        if remove_bucket {
376            owned_scopes.remove(&key);
377        }
378        shrink_map_if_sparse(&mut owned_scopes, Self::MIN_RETAINED_SCOPE_CAPACITY);
379        removed
380    }
381
382    fn partition_scopes(
383        &self,
384        should_remove: impl Fn(&ScopeEntry) -> bool,
385    ) -> Vec<Rc<RefCell<ScopeEntry>>> {
386        let mut owned_scopes = self.owned_scopes.borrow_mut();
387        let mut retained = HashMap::default();
388        let mut removed = Vec::new();
389        for (key, mut bucket) in owned_scopes.drain() {
390            let mut retained_bucket = OwnedScopeBucket::new();
391            for entry in bucket.drain(..) {
392                if should_remove(&entry.borrow()) {
393                    removed.push(entry);
394                } else {
395                    retained_bucket.push(entry);
396                }
397            }
398            if !retained_bucket.is_empty() {
399                retained.insert(key, retained_bucket);
400            }
401        }
402        *owned_scopes = retained;
403        shrink_map_if_sparse(&mut owned_scopes, Self::MIN_RETAINED_SCOPE_CAPACITY);
404        removed
405    }
406
407    fn debug_stats(&self) -> SnapshotStateObserverDebugStats {
408        let owned_scopes = self.owned_scopes.borrow();
409        let indexed_scopes = self.indexed_scopes.borrow();
410        let owned_scope_cap =
411            owned_scopes.capacity() + owned_scopes.values().map(SmallVec::capacity).sum::<usize>();
412        let scopes_cap = owned_scope_cap + indexed_scopes.capacity();
413        let mut observed_state_count = 0;
414        let mut observed_state_capacity = 0;
415        let mut stateless_scope_count = 0;
416
417        for entry in indexed_scopes.values() {
418            let entry = entry.borrow();
419            observed_state_count += entry.observed.len();
420            observed_state_capacity += entry.observed.capacity();
421            stateless_scope_count += usize::from(entry.observed.is_empty());
422        }
423
424        SnapshotStateObserverDebugStats {
425            scopes_len: indexed_scopes.len(),
426            scopes_cap,
427            stateless_scope_count,
428            observed_state_count,
429            observed_state_capacity,
430        }
431    }
432
433    fn run_with_read_observer<R>(&self, block: impl FnOnce() -> R) -> R {
434        use crate::snapshot_v2::{
435            current_snapshot_reads_into, take_transparent_observer_mutable_snapshot_reusing,
436        };
437
438        if current_snapshot_reads_into(&self.read_dispatcher) {
439            return block();
440        }
441
442        let mut snapshot = take_transparent_observer_mutable_snapshot_reusing(
443            Some(self.read_dispatcher.clone()),
444            None,
445            self.read_snapshot.take(),
446        );
447        let result = snapshot.enter(block);
448        snapshot.dispose();
449        if Arc::get_mut(&mut snapshot).is_some() && !snapshot.has_pending_changes() {
450            self.read_snapshot.replace(Some(snapshot));
451        }
452        result
453    }
454
455    fn handle_apply(&self, modified: &[Arc<dyn StateObject>]) {
456        if modified.is_empty() {
457            return;
458        }
459
460        let mut seen_scope_ids: HashSet<usize> = HashSet::default();
461        let mut to_notify: Vec<Rc<RefCell<ScopeEntry>>> = Vec::new();
462        {
463            let observed_to_scopes = self.observed_to_scopes.borrow();
464            let indexed_scopes = self.indexed_scopes.borrow();
465            for state in modified {
466                if let Some(scope_ids) = observed_to_scopes.get(&state.object_id().as_usize()) {
467                    let mut ordered_scope_ids: SmallVec<[usize; 8]> =
468                        scope_ids.iter().copied().collect();
469                    ordered_scope_ids.sort_unstable();
470                    for scope_id in ordered_scope_ids {
471                        if seen_scope_ids.insert(scope_id)
472                            && let Some(entry) = indexed_scopes.get(&scope_id)
473                        {
474                            to_notify.push(entry.clone());
475                        }
476                    }
477                }
478            }
479        }
480
481        if to_notify.is_empty() {
482            return;
483        }
484
485        for entry in to_notify {
486            let executor = self.executor.clone();
487            executor(Box::new(move || {
488                if let Ok(entry) = entry.try_borrow() {
489                    entry.notify();
490                }
491            }));
492        }
493    }
494
495    fn replace_observed_ids(&self, entry: &Rc<RefCell<ScopeEntry>>, collected: &mut ObservedIds) {
496        let (entry_id, previous) = {
497            let mut entry_mut = entry.borrow_mut();
498            if entry_mut.observed.iter().eq(collected.iter()) {
499                entry_mut.observed.take_leases(collected);
500                collected.clear();
501                return;
502            }
503            let entry_id = entry_mut.id;
504            let previous = std::mem::replace(&mut entry_mut.observed, collected.take_sized());
505            (entry_id, previous)
506        };
507        let entry_ref = entry.borrow();
508        self.unregister_observed_ids(entry_id, &previous);
509        self.register_observed_ids(entry_id, &entry_ref.observed);
510    }
511
512    fn register_observed_ids(&self, entry_id: usize, observed: &ObservedIds) {
513        let mut observed_to_scopes = self.observed_to_scopes.borrow_mut();
514        for state_id in observed.iter() {
515            let scope_ids = observed_to_scopes.entry(state_id).or_default();
516            scope_ids.insert(entry_id);
517        }
518    }
519
520    fn unregister_observed_ids(&self, entry_id: usize, observed: &ObservedIds) {
521        let mut observed_to_scopes = self.observed_to_scopes.borrow_mut();
522        let mut emptied = SmallVec::<[StateObjectId; MAX_OBSERVED_STATES]>::new();
523        for state_id in observed.iter() {
524            if let Some(scope_ids) = observed_to_scopes.get_mut(&state_id) {
525                scope_ids.remove(&entry_id);
526                if scope_ids.is_empty() {
527                    emptied.push(state_id);
528                }
529            }
530        }
531        for state_id in emptied {
532            observed_to_scopes.remove(&state_id);
533        }
534        shrink_map_if_sparse(&mut observed_to_scopes, Self::MIN_RETAINED_SCOPE_CAPACITY);
535    }
536
537    fn unregister_entry(&self, entry: &Rc<RefCell<ScopeEntry>>) {
538        let (entry_id, observed) = {
539            let mut entry_mut = entry.borrow_mut();
540            let observed = std::mem::replace(&mut entry_mut.observed, ObservedIds::new());
541            (entry_mut.id, observed)
542        };
543        self.unregister_observed_ids(entry_id, &observed);
544        self.indexed_scopes.borrow_mut().remove(&entry_id);
545    }
546}
547
548fn shrink_map_if_sparse<K, V>(map: &mut HashMap<K, V>, min_retained_capacity: usize)
549where
550    K: Eq + std::hash::Hash,
551{
552    if map.capacity() <= map.len().max(min_retained_capacity).saturating_mul(4) {
553        return;
554    }
555
556    let retained = map.len().max(min_retained_capacity);
557    let mut rebuilt = HashMap::default();
558    rebuilt.reserve(retained);
559    rebuilt.extend(map.drain());
560    *map = rebuilt;
561}
562
563#[derive(Default)]
564struct ReadObservationStack {
565    targets: Vec<Rc<RefCell<ObservedIds>>>,
566    depth: usize,
567}
568
569impl ReadObservationStack {
570    fn push(&mut self) -> Rc<RefCell<ObservedIds>> {
571        if self.depth == self.targets.len() {
572            self.targets.push(Rc::new(RefCell::new(ObservedIds::new())));
573        }
574        let target = Rc::clone(&self.targets[self.depth]);
575        self.depth += 1;
576        target
577    }
578
579    fn last(&self) -> Option<&Rc<RefCell<ObservedIds>>> {
580        self.depth.checked_sub(1).map(|index| &self.targets[index])
581    }
582
583    fn pop(&mut self) -> Rc<RefCell<ObservedIds>> {
584        self.depth -= 1;
585        Rc::clone(&self.targets[self.depth])
586    }
587}
588
589enum ObservedIds {
590    Small(SmallVec<[ObservedState; 1]>),
591    Large(Box<HashMap<StateObjectId, Option<Rc<dyn Any>>>>),
592}
593
594struct ObservedState {
595    id: StateObjectId,
596    _lease: Option<Rc<dyn Any>>,
597}
598
599impl ObservedIds {
600    fn new() -> Self {
601        ObservedIds::Small(SmallVec::new())
602    }
603
604    fn insert(&mut self, state: &dyn StateObject) {
605        let id = state.object_id().as_usize();
606        match self {
607            ObservedIds::Small(small) => {
608                if small.iter().any(|observed| observed.id == id) {
609                    return;
610                }
611                if small.len() < MAX_OBSERVED_STATES {
612                    small.push(ObservedState {
613                        id,
614                        _lease: state.observation_lease(),
615                    });
616                } else {
617                    let mut large =
618                        HashMap::with_capacity_and_hasher(small.len() + 1, Default::default());
619                    for observed in small.drain(..) {
620                        large.insert(observed.id, observed._lease);
621                    }
622                    large.insert(id, state.observation_lease());
623                    *self = ObservedIds::Large(Box::new(large));
624                }
625            }
626            ObservedIds::Large(large) => {
627                large.entry(id).or_insert_with(|| state.observation_lease());
628            }
629        }
630    }
631
632    fn clear(&mut self) {
633        match self {
634            ObservedIds::Small(small) => small.clear(),
635            ObservedIds::Large(_) => *self = ObservedIds::new(),
636        }
637    }
638
639    /// Keeps `fresh`'s observation leases for the states both name, leaving
640    /// the old ones in `fresh` to drop. An id is a state's address, which a
641    /// new state can take after the old one is dropped, and only its own
642    /// lease keeps the new one observed.
643    fn take_leases(&mut self, fresh: &mut ObservedIds) {
644        match (self, fresh) {
645            (ObservedIds::Small(kept), ObservedIds::Small(fresh)) => {
646                for (kept, fresh) in kept.iter_mut().zip(fresh.iter_mut()) {
647                    std::mem::swap(&mut kept._lease, &mut fresh._lease);
648                }
649            }
650            (kept, fresh) => std::mem::swap(kept, fresh),
651        }
652    }
653
654    fn take_sized(&mut self) -> ObservedIds {
655        match self {
656            ObservedIds::Small(small) => ObservedIds::Small(small.drain(..).collect()),
657            ObservedIds::Large(_) => std::mem::replace(self, ObservedIds::new()),
658        }
659    }
660
661    fn is_empty(&self) -> bool {
662        match self {
663            ObservedIds::Small(small) => small.is_empty(),
664            ObservedIds::Large(large) => large.is_empty(),
665        }
666    }
667
668    fn len(&self) -> usize {
669        match self {
670            ObservedIds::Small(small) => small.len(),
671            ObservedIds::Large(large) => large.len(),
672        }
673    }
674
675    fn capacity(&self) -> usize {
676        match self {
677            ObservedIds::Small(small) => small.capacity(),
678            ObservedIds::Large(large) => large.capacity(),
679        }
680    }
681
682    fn iter(&self) -> impl Iterator<Item = StateObjectId> + '_ {
683        let (small, large) = match self {
684            ObservedIds::Small(small) => (Some(small.as_slice()), None),
685            ObservedIds::Large(large) => (None, Some(large)),
686        };
687        small
688            .into_iter()
689            .flatten()
690            .map(|observed| observed.id)
691            .chain(large.into_iter().flat_map(|states| states.keys().copied()))
692    }
693}
694
695const MAX_OBSERVED_STATES: usize = 8;
696
697struct ScopeEntry {
698    id: usize,
699    scope: Box<dyn Any>,
700    on_changed: Rc<dyn ScopeChangedCallback>,
701    observed: ObservedIds,
702}
703
704impl ScopeEntry {
705    fn new<T>(id: usize, scope: T, on_changed: Rc<dyn ScopeChangedCallback>) -> Self
706    where
707        T: Any + 'static,
708    {
709        Self {
710            id,
711            scope: Box::new(scope),
712            on_changed,
713            observed: ObservedIds::new(),
714        }
715    }
716
717    fn callback_reusing<F: Fn(&dyn Any) + 'static>(
718        &mut self,
719        callback: F,
720    ) -> Rc<dyn ScopeChangedCallback> {
721        if let Some(stored) = Rc::get_mut(&mut self.on_changed)
722            .and_then(|stored| (stored as &mut dyn Any).downcast_mut::<F>())
723        {
724            *stored = callback;
725            Rc::clone(&self.on_changed)
726        } else {
727            Rc::new(callback)
728        }
729    }
730
731    fn update<T>(&mut self, new_scope: &T, on_changed: Rc<dyn ScopeChangedCallback>)
732    where
733        T: Any + Clone + 'static,
734    {
735        self.update_scope(new_scope);
736        self.on_changed = on_changed;
737    }
738
739    fn update_scope<T>(&mut self, new_scope: &T)
740    where
741        T: Any + Clone + 'static,
742    {
743        match self.scope.downcast_mut::<T>() {
744            Some(stored) => stored.clone_from(new_scope),
745            None => self.scope = Box::new(new_scope.clone()),
746        }
747    }
748
749    fn matches_scope<T>(&self, scope: &T) -> bool
750    where
751        T: Any + Eq + 'static,
752    {
753        self.scope
754            .downcast_ref::<T>()
755            .is_some_and(|stored| stored == scope)
756    }
757
758    fn notify(&self) {
759        (self.on_changed)(self.scope.as_ref());
760    }
761}
762
763#[cfg(test)]
764#[path = "tests/snapshot_state_observer_tests.rs"]
765mod tests;