Skip to main content

guinea_core/actor/
addr.rs

1use crate::actor::cancel::Cancel;
2use crate::actor::envelope::{Envelope, FnEnvelope, MessageEnvelope};
3use crate::actor::event_bus::builder::EventSubscription;
4use crate::actor::event_bus::subscribe::{BusSubscription, Event};
5use crate::actor::event_bus::{EventBus, GlobalEventBus};
6use crate::actor::shape::name;
7use crate::actor::traits::Handler;
8use crate::actor::{Cx, UiThreadToken};
9use crate::actor::{ManagedActor, short_type_name};
10use crate::lifecycle_tracker::LifecycleTracker;
11use crate::scope::Scope;
12use crate::trace::{self, Bus, Cause, Point};
13use std::any::Any;
14use std::cell::{Cell, RefCell};
15use std::collections::{HashMap, VecDeque};
16use std::marker::PhantomData;
17use std::rc::{Rc, Weak};
18use std::sync::atomic::{AtomicUsize, Ordering};
19
20static NEXT_ID: AtomicUsize = AtomicUsize::new(1);
21
22thread_local! {
23    pub static REGISTRY: RefCell<HashMap<usize, Box<dyn Any>>> = RefCell::new(HashMap::new());
24}
25
26/// The actor registered under `id`, if it is still there and is an `A`.
27///
28/// A clone, taken with the registry borrowed only for as long as it takes to
29/// clone: whatever the caller does with it next - send, and so run handlers
30/// that create or dispose actors - finds the registry free.
31pub(crate) fn registered<A: 'static>(id: usize) -> Option<Addr<A>> {
32    REGISTRY.with(|reg| {
33        reg.borrow()
34            .get(&id)
35            .and_then(|addr| addr.downcast_ref::<Addr<A>>())
36            .cloned()
37    })
38}
39
40pub struct Addr<A: 'static> {
41    pub(super) id: usize,
42    pub(super) guard: UiThreadToken,
43    state: Rc<RefCell<A>>,
44    queue: Rc<RefCell<VecDeque<Box<dyn Envelope<A>>>>>,
45    is_processing: Rc<Cell<bool>>,
46    counter: Rc<&'static str>,
47    cancel: Cancel,
48    /// What it hears, for as long as it lives: disposing it ends them.
49    subscriptions: Rc<RefCell<Vec<BusSubscription>>>,
50    /// The scope it belongs to and its window's bus, when it has them.
51    home: Rc<RefCell<Option<Home>>>,
52}
53
54struct Home {
55    scope: Scope,
56    bus: Weak<EventBus>,
57}
58
59impl<A: 'static> Clone for Addr<A> {
60    fn clone(&self) -> Self {
61        Self {
62            id: self.id,
63            state: self.state.clone(),
64            guard: self.guard.clone(),
65            queue: self.queue.clone(),
66            is_processing: self.is_processing.clone(),
67            counter: self.counter.clone(),
68            cancel: self.cancel.clone(),
69            subscriptions: self.subscriptions.clone(),
70            home: self.home.clone(),
71        }
72    }
73}
74
75/// Keeps what an actor subscribed to on the actor itself.
76struct Held(Rc<RefCell<Vec<BusSubscription>>>);
77
78impl LifecycleTracker for Held {
79    fn track_loop<T: 'static>(&self, _handle: T) {}
80    fn track_actor<A: 'static>(&self, _addr: &Addr<A>) {}
81
82    fn track_sub(&self, subscription: BusSubscription) {
83        self.0.borrow_mut().push(subscription);
84    }
85}
86
87impl<A: 'static> Addr<A> {
88    pub fn new_managed(state: A, token: UiThreadToken, tracker: &impl LifecycleTracker) -> Self
89    where
90        A: ManagedActor,
91    {
92        let addr = Self::new(state, token, tracker);
93
94        A::Bus::subscribe_into(addr.clone(), tracker);
95
96        addr
97    }
98
99    /// An actor a scope owns. What its manifest subscribes to is held by the
100    /// actor, and ends when the scope disposes it.
101    pub fn new_managed_scoped(state: A, token: UiThreadToken) -> Self
102    where
103        A: ManagedActor,
104    {
105        let addr = Self::new(state, token, &crate::lifecycle_tracker::NullTracker);
106        A::Bus::subscribe_into(addr.clone(), &Held(addr.subscriptions.clone()));
107        addr
108    }
109
110    /// Where the actor lives: the scope that owns it, and that scope's
111    /// window bus when it is in a window. What `subscribe_on` reaches the
112    /// window bus through, and notes a listener on.
113    #[doc(hidden)]
114    pub fn live_in(&self, scope: Scope, bus: Option<&Rc<EventBus>>) {
115        *self.home.borrow_mut() = Some(Home {
116            scope,
117            bus: bus.map(Rc::downgrade).unwrap_or_default(),
118        });
119    }
120
121    /// Hears `M` on `bus` for as long as the actor lives: disposing it ends
122    /// the subscription.
123    pub fn subscribe_on<M: Event>(&self, bus: Bus)
124    where
125        A: Handler<M>,
126    {
127        let (scope, window) = match self.home.borrow().as_ref() {
128            Some(home) => (Some(home.scope), home.bus.upgrade()),
129            None => (None, None),
130        };
131
132        let on = match bus {
133            Bus::Global => GlobalEventBus::bus(),
134            Bus::Window => window.unwrap_or_else(|| {
135                panic!(
136                    "{} lives in no window, so there is no window bus to hear {} on",
137                    short_type_name::<A>(),
138                    short_type_name::<M>()
139                )
140            }),
141        };
142
143        if let Some(scope) = scope {
144            scope.note_listener(name::<M>(), Some(name::<A>()), bus);
145        }
146
147        let subscription = on.subscribe::<A, M>(self.clone());
148        self.subscriptions.borrow_mut().push(subscription);
149    }
150
151    /// Whether the scope it lives in is asleep, or gone: what a bus carries
152    /// is not for it.
153    pub(crate) fn is_asleep(&self) -> bool {
154        self.home
155            .borrow()
156            .as_ref()
157            .is_some_and(|home| !home.scope.is_awake())
158    }
159
160    pub fn new_scoped(state: A, token: UiThreadToken) -> Self {
161        Self::new(state, token, &crate::lifecycle_tracker::NullTracker)
162    }
163
164    pub fn new(state: A, guard: UiThreadToken, tracker: &impl LifecycleTracker) -> Self {
165        let id = NEXT_ID.fetch_add(1, Ordering::Relaxed);
166        let addr = Self {
167            id,
168            guard,
169            state: Rc::new(RefCell::new(state)),
170            queue: Rc::new(RefCell::new(VecDeque::new())),
171            is_processing: Rc::new(Cell::new(false)),
172            counter: Rc::new(short_type_name::<A>()),
173            cancel: Cancel::new(),
174            subscriptions: Rc::new(RefCell::new(Vec::new())),
175            home: Rc::new(RefCell::new(None)),
176        };
177
178        let addr_clone = addr.clone();
179        REGISTRY.with(|reg| {
180            reg.borrow_mut().insert(id, Box::new(addr_clone));
181        });
182
183        tracker.track_actor(&addr);
184        addr
185    }
186
187    pub fn apply<F>(&self, f: F)
188    where
189        F: FnOnce(&mut A, &Cx<A>) + Send + 'static,
190    {
191        self.queue.borrow_mut().push_back(Box::new(FnEnvelope {
192            func: Some(f),
193            cause: trace::current(),
194            phantom: PhantomData,
195        }));
196
197        self.process_queue();
198    }
199
200    pub fn handler<M>(&self, msg: M) -> impl Fn() + 'static
201    where
202        M: Clone + 'static,
203        A: Handler<M>,
204    {
205        let addr = self.clone();
206        move || addr.do_send(msg.clone())
207    }
208
209    pub fn handler_with<M, T, F>(&self, f: F) -> impl Fn(T) + 'static
210    where
211        F: Fn(T) -> M + 'static,
212        M: 'static,
213        A: Handler<M>,
214    {
215        let addr = self.clone();
216        move |arg| addr.do_send(f(arg))
217    }
218
219    pub fn handler_with2<M, T1, T2, F>(&self, f: F) -> impl Fn(T1, T2) + 'static
220    where
221        F: Fn(T1, T2) -> M + 'static,
222        M: 'static,
223        A: Handler<M>,
224    {
225        let addr = self.clone();
226        move |arg1, arg2| addr.do_send(f(arg1, arg2))
227    }
228
229    pub fn send<M>(&self, msg: M)
230    where
231        M: 'static,
232        A: Handler<M>,
233    {
234        self.do_send(msg);
235    }
236
237    #[cfg(feature = "test-utils")]
238    pub fn send_test<M>(&self, msg: M) -> crate::test_kit::Interaction<()>
239    where
240        M: 'static,
241        A: Handler<M>,
242    {
243        self.do_send(msg);
244        crate::test_kit::Interaction::new(())
245    }
246
247    fn do_send<M>(&self, msg: M)
248    where
249        M: 'static,
250        A: Handler<M>,
251    {
252        self.send_under(msg, trace::current());
253    }
254
255    /// Queues `msg` as caused by `parent`: for a message whose cause crossed a
256    /// thread or a background task to get here.
257    pub(crate) fn send_under<M>(&self, msg: M, parent: Option<Cause>)
258    where
259        M: 'static,
260        A: Handler<M>,
261    {
262        let cause = trace::mark_under(parent, || Point::Send {
263            actor: short_type_name::<A>(),
264            message: short_type_name::<M>(),
265        });
266
267        self.queue.borrow_mut().push_back(Box::new(MessageEnvelope {
268            message: Some(msg),
269            cause,
270        }));
271
272        self.process_queue();
273    }
274
275    pub fn get_token(&self) -> UiThreadToken {
276        self.guard.clone()
277    }
278    pub fn strong_count_ptr(&self) -> Rc<&'static str> {
279        self.counter.clone()
280    }
281
282    pub fn id(&self) -> usize {
283        self.id
284    }
285
286    /// The token every task this actor spawned is guarded by; cancelled by
287    /// [`Addr::dispose`], and so by the teardown that owns the actor.
288    pub fn cancellation(&self) -> Cancel {
289        self.cancel.clone()
290    }
291
292    pub fn debug_snapshot(&self) -> String
293    where
294        A: std::fmt::Debug,
295    {
296        match self.state.try_borrow() {
297            Ok(state) => format!("{:#?}", *state),
298            Err(_) => "<handling a message>".to_string(),
299        }
300    }
301
302    /// Takes the actor out of the registry and ends its background work: what
303    /// it spawned is dropped where it last awaited, instead of running on with
304    /// nowhere to answer.
305    pub fn dispose(&self) {
306        self.cancel.cancel();
307        self.subscriptions.borrow_mut().clear();
308
309        let gone = REGISTRY.with(|reg| reg.borrow_mut().remove(&self.id));
310        drop(gone);
311    }
312
313    fn process_queue(&self) {
314        if self.is_processing.get() {
315            return;
316        }
317        self.is_processing.set(true);
318
319        crate::notify::turn(|| self.drain_queue());
320    }
321
322    fn drain_queue(&self) {
323        loop {
324            let mut envelope = {
325                let mut q = self.queue.borrow_mut();
326                match q.pop_front() {
327                    Some(e) => e,
328                    None => {
329                        self.is_processing.set(false);
330                        break;
331                    }
332                }
333            };
334
335            {
336                let mut state_guard = self.state.borrow_mut();
337                Envelope::<A>::handle(envelope.as_mut(), &mut *state_guard, self);
338            }
339            crate::devtools::changed(|| crate::devtools::Change::ActorHandled { id: self.id });
340        }
341
342        self.is_processing.set(false);
343    }
344}
345
346#[cfg(test)]
347mod tests {
348    use super::*;
349
350    #[test]
351    fn an_actor_whose_scope_was_removed_hears_nothing_more() {
352        let scope = crate::scope::ScopeTree::new();
353        let addr = Addr::new_scoped((), UiThreadToken::dangerously_create_token_unchecked());
354        addr.live_in(scope.scope(), Some(&Rc::new(EventBus::new())));
355
356        scope.remove();
357
358        assert!(addr.is_asleep(), "a removed scope read as awake");
359    }
360}