Skip to main content

guinea_core/actor/
addr.rs

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