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
23pub(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 subscriptions: Rc<RefCell<Vec<BusSubscription>>>,
47 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 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 #[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 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 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 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 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 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}