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
26pub(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 subscriptions: Rc<RefCell<Vec<BusSubscription>>>,
50 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
75struct 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 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 #[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 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 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 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 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 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}