reifydb_runtime/actor/system/native/
mod.rs1#![allow(clippy::disallowed_types)]
5
6mod pool;
7
8use std::{
9 any::Any,
10 error, fmt,
11 fmt::{Debug, Formatter},
12 mem,
13 sync::{Arc, Weak},
14 time,
15 time::Duration,
16};
17
18use crossbeam_channel::{Receiver, RecvTimeoutError as CcRecvTimeoutError};
19
20use crate::{
21 actor::{
22 context::CancellationToken, system::native::pool::PoolActorHandle, timers::scheduler::SchedulerHandle,
23 traits::Actor,
24 },
25 context::clock::Clock,
26 pool::Pools,
27 sync::mutex::Mutex,
28};
29
30struct ActorSystemInner {
31 cancel: CancellationToken,
32 scheduler: SchedulerHandle,
33 clock: Clock,
34 pools: Pools,
35 wakers: Mutex<Vec<Arc<dyn Fn() + Send + Sync>>>,
36 keepalive: Mutex<Vec<Box<dyn Any + Send + Sync>>>,
37 done_rxs: Mutex<Vec<Receiver<()>>>,
38 children: Mutex<Vec<ActorSystem>>,
39}
40
41#[derive(Clone)]
42pub struct ActorSystem {
43 inner: Arc<ActorSystemInner>,
44}
45
46impl ActorSystem {
47 pub fn new(pools: Pools, clock: Clock) -> Self {
48 let scheduler = SchedulerHandle::new(pools.system_pool().clone());
49
50 Self {
51 inner: Arc::new(ActorSystemInner {
52 cancel: CancellationToken::new(),
53 scheduler,
54 clock,
55 pools,
56 wakers: Mutex::new(Vec::new()),
57 keepalive: Mutex::new(Vec::new()),
58 done_rxs: Mutex::new(Vec::new()),
59 children: Mutex::new(Vec::new()),
60 }),
61 }
62 }
63
64 pub fn scope(&self) -> Self {
65 let child = Self {
66 inner: Arc::new(ActorSystemInner {
67 cancel: self.inner.cancel.child_token(),
68 scheduler: self.inner.scheduler.shared(),
69 clock: self.inner.clock.clone(),
70 pools: self.inner.pools.clone(),
71 wakers: Mutex::new(Vec::new()),
72 keepalive: Mutex::new(Vec::new()),
73 done_rxs: Mutex::new(Vec::new()),
74 children: Mutex::new(Vec::new()),
75 }),
76 };
77 self.inner.children.lock().push(child.clone());
78 child
79 }
80
81 pub fn pools(&self) -> Pools {
82 self.inner.pools.clone()
83 }
84
85 pub fn spawner(&self) -> ActorSpawner {
86 ActorSpawner {
87 inner: Arc::downgrade(&self.inner),
88 clock: self.inner.clock.clone(),
89 }
90 }
91
92 pub fn cancellation_token(&self) -> CancellationToken {
93 self.inner.cancel.clone()
94 }
95
96 pub fn is_cancelled(&self) -> bool {
97 self.inner.cancel.is_cancelled()
98 }
99
100 pub fn shutdown(&self) {
101 self.inner.cancel.cancel();
102
103 {
104 let mut children = self.inner.children.lock();
105 for child in children.iter() {
106 child.shutdown();
107 }
108 children.clear();
109 }
110
111 let wakers = mem::take(&mut *self.inner.wakers.lock());
112 for waker in &wakers {
113 waker();
114 }
115 drop(wakers);
116
117 self.inner.keepalive.lock().clear();
118 }
119
120 pub(crate) fn register_waker(&self, f: Arc<dyn Fn() + Send + Sync>) {
121 self.inner.wakers.lock().push(f);
122 }
123
124 pub(crate) fn register_keepalive(&self, cell: Box<dyn Any + Send + Sync>) {
125 self.inner.keepalive.lock().push(cell);
126 }
127
128 pub(crate) fn register_done_rx(&self, rx: Receiver<()>) {
129 self.inner.done_rxs.lock().push(rx);
130 }
131
132 pub fn join(&self) -> Result<(), JoinError> {
133 self.join_timeout(Duration::from_secs(5))
134 }
135
136 #[allow(clippy::disallowed_methods)]
137 pub fn join_timeout(&self, timeout: Duration) -> Result<(), JoinError> {
138 let deadline = time::Instant::now() + timeout;
139 let rxs: Vec<_> = mem::take(&mut *self.inner.done_rxs.lock());
140 for rx in rxs {
141 let remaining = deadline.saturating_duration_since(time::Instant::now());
142 match rx.recv_timeout(remaining) {
143 Ok(()) => {}
144 Err(CcRecvTimeoutError::Disconnected) => {}
145 Err(CcRecvTimeoutError::Timeout) => {
146 return Err(JoinError::new("timed out waiting for actors to stop"));
147 }
148 }
149 }
150 Ok(())
151 }
152
153 pub fn scheduler(&self) -> &SchedulerHandle {
154 &self.inner.scheduler
155 }
156
157 pub fn clock(&self) -> &Clock {
158 &self.inner.clock
159 }
160
161 pub fn spawn_system<A: Actor>(&self, name: &str, actor: A) -> ActorHandle<A::Message>
162 where
163 A::State: Send,
164 {
165 pool::spawn_on_pool(self, name, actor, self.inner.pools.system_pool())
166 }
167
168 pub fn spawn_query<A: Actor>(&self, name: &str, actor: A) -> ActorHandle<A::Message>
169 where
170 A::State: Send,
171 {
172 pool::spawn_on_pool(self, name, actor, self.inner.pools.query_pool())
173 }
174
175 pub fn spawn_commit<A: Actor>(&self, name: &str, actor: A) -> ActorHandle<A::Message>
176 where
177 A::State: Send,
178 {
179 pool::spawn_on_pool(self, name, actor, self.inner.pools.commit_pool())
180 }
181
182 pub fn spawn_background<A: Actor>(&self, name: &str, actor: A) -> ActorHandle<A::Message>
183 where
184 A::State: Send,
185 {
186 pool::spawn_on_pool(self, name, actor, self.inner.pools.background_pool())
187 }
188}
189
190impl Debug for ActorSystem {
191 fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
192 f.debug_struct("ActorSystem").field("cancelled", &self.is_cancelled()).finish_non_exhaustive()
193 }
194}
195
196#[derive(Clone)]
197pub struct ActorSpawner {
198 inner: Weak<ActorSystemInner>,
199 clock: Clock,
200}
201
202impl ActorSpawner {
203 fn system(&self) -> ActorSystem {
204 ActorSystem {
205 inner: self.inner.upgrade().expect("runtime already shut down: cannot spawn actor"),
206 }
207 }
208
209 pub fn clock(&self) -> &Clock {
210 &self.clock
211 }
212
213 pub fn pools(&self) -> Pools {
214 self.system().pools()
215 }
216
217 pub fn is_alive(&self) -> bool {
218 self.inner.strong_count() > 0
219 }
220
221 pub fn cancellation_token(&self) -> Option<CancellationToken> {
222 self.inner.upgrade().map(|inner| inner.cancel.clone())
223 }
224
225 pub fn scope(&self) -> ActorSpawner {
226 self.system().scope().spawner()
227 }
228
229 pub fn shutdown(&self) {
230 if let Some(inner) = self.inner.upgrade() {
231 ActorSystem {
232 inner,
233 }
234 .shutdown();
235 }
236 }
237
238 pub fn spawn_system<A: Actor>(&self, name: &str, actor: A) -> ActorHandle<A::Message>
239 where
240 A::State: Send,
241 {
242 self.system().spawn_system(name, actor)
243 }
244
245 pub fn spawn_query<A: Actor>(&self, name: &str, actor: A) -> ActorHandle<A::Message>
246 where
247 A::State: Send,
248 {
249 self.system().spawn_query(name, actor)
250 }
251
252 pub fn spawn_commit<A: Actor>(&self, name: &str, actor: A) -> ActorHandle<A::Message>
253 where
254 A::State: Send,
255 {
256 self.system().spawn_commit(name, actor)
257 }
258
259 pub fn spawn_background<A: Actor>(&self, name: &str, actor: A) -> ActorHandle<A::Message>
260 where
261 A::State: Send,
262 {
263 self.system().spawn_background(name, actor)
264 }
265}
266
267impl Debug for ActorSpawner {
268 fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
269 f.debug_struct("ActorSpawner").field("alive", &self.is_alive()).finish_non_exhaustive()
270 }
271}
272
273pub type ActorHandle<M> = PoolActorHandle<M>;
274
275#[derive(Debug)]
276pub struct JoinError {
277 message: String,
278}
279
280impl JoinError {
281 pub fn new(message: impl Into<String>) -> Self {
282 Self {
283 message: message.into(),
284 }
285 }
286}
287
288impl fmt::Display for JoinError {
289 fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
290 write!(f, "actor join failed: {}", self.message)
291 }
292}
293
294impl error::Error for JoinError {}
295
296#[cfg(test)]
297mod tests {
298 use std::sync;
299
300 use super::*;
301 use crate::{
302 actor::{context::Context, traits::Directive},
303 pool::{PoolConfig, Pools},
304 };
305
306 fn test_system() -> ActorSystem {
307 let pools = Pools::new(PoolConfig::default());
308 ActorSystem::new(pools, Clock::Real)
309 }
310
311 struct CounterActor;
312
313 #[derive(Debug)]
314 enum CounterMessage {
315 Inc,
316 Get(sync::mpsc::Sender<i64>),
317 Stop,
318 }
319
320 impl Actor for CounterActor {
321 type State = i64;
322 type Message = CounterMessage;
323
324 fn init(&self, _ctx: &Context<Self::Message>) -> Self::State {
325 0
326 }
327
328 fn handle(
329 &self,
330 state: &mut Self::State,
331 msg: Self::Message,
332 _ctx: &Context<Self::Message>,
333 ) -> Directive {
334 match msg {
335 CounterMessage::Inc => *state += 1,
336 CounterMessage::Get(tx) => {
337 let _ = tx.send(*state);
338 }
339 CounterMessage::Stop => return Directive::Stop,
340 }
341 Directive::Continue
342 }
343 }
344
345 #[test]
346 fn test_spawn_and_send() {
347 let system = test_system();
348 let handle = system.spawn_system("counter", CounterActor);
349
350 let actor_ref = handle.actor_ref().clone();
351 actor_ref.send(CounterMessage::Inc).unwrap();
352 actor_ref.send(CounterMessage::Inc).unwrap();
353 actor_ref.send(CounterMessage::Inc).unwrap();
354
355 let (tx, rx) = sync::mpsc::channel();
356 actor_ref.send(CounterMessage::Get(tx)).unwrap();
357
358 let value = rx.recv().unwrap();
359 assert_eq!(value, 3);
360
361 actor_ref.send(CounterMessage::Stop).unwrap();
362 handle.join().unwrap();
363 }
364
365 #[test]
366 fn test_shutdown_join() {
367 let system = test_system();
368
369 for i in 0..5 {
371 system.spawn_system(&format!("counter-{i}"), CounterActor);
372 }
373
374 system.shutdown();
376 system.join().unwrap();
377 }
378}