Skip to main content

reifydb_runtime/actor/system/native/
mod.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4#![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		// Spawn several actors
370		for i in 0..5 {
371			system.spawn_system(&format!("counter-{i}"), CounterActor);
372		}
373
374		// Shutdown cancels all actors; join waits for them to finish
375		system.shutdown();
376		system.join().unwrap();
377	}
378}