actor_framework/actor.rs
1//! # Generic Actor Server
2//!
3//! This module defines the `ResourceActor`, the core component that manages the lifecycle
4//! and state of entities. It implements the "Server" side of the Actor Model, processing
5//! messages sequentially and ensuring exclusive access to the entity store.
6
7use crate::client::ResourceClient;
8use crate::entity::ActorEntity;
9use crate::error::FrameworkError;
10use crate::message::ResourceRequest;
11use std::collections::HashMap;
12use tokio::sync::mpsc;
13use tracing::{debug, info, warn};
14
15/// The generic actor that manages a collection of entities.
16///
17/// # Architecture Note
18/// This struct is the "Server" half of the actor. It owns the state (`store`) and
19/// the receiver end of the channel.
20///
21/// **Concurrency Model**:
22/// Even though we might have 1000 `ResourceActor` instances running, each one
23/// processes its own messages *sequentially* in a loop. This means we don't need
24/// `Mutex` or `RwLock` for the `store`! The "Actor Model" gives us safety through
25/// exclusive ownership of state within the task.
26/// ## ResourceActor
27///
28/// The `ResourceActor<T>` struct is the *server* side of the framework. It owns the in‑memory store for a given entity type `T: ActorEntity` and processes all incoming `ResourceRequest<T>` messages sequentially. Each actor runs in its own Tokio task, guaranteeing exclusive access to its state without any locking.
29///
30/// * **Concurrency model** – each actor processes one message at a time, eliminating data races.
31/// * **Context injection** – a user‑provided `Context` is passed to every lifecycle hook, enabling dependency injection.
32/// * **Uniform API** – works with any entity that implements `ActorEntity`, providing a generic CRUD + Action implementation.
33///
34/// # Usage Pattern
35///
36/// The canonical way to create and wire actors is:
37///
38/// 1. **Create**: Call `ResourceActor::new()` to get the `actor` (server) and `client` (interface).
39/// 2. **Wire**: Pass dependencies (other clients) into `actor.run(context)`.
40/// 3. **Run**: Spawn the actor's run loop in a background task.
41///
42/// ```rust
43/// use actor_framework::{ActorEntity, ResourceActor};
44/// use async_trait::async_trait;
45///
46/// // Minimal Entity Definition
47/// #[derive(Clone, Debug)] struct MyEntity { id: u32 }
48/// #[derive(Debug)] struct MyCreate;
49/// #[derive(Debug)] struct MyUpdate;
50/// #[derive(Debug)] enum MyAction {}
51/// #[derive(Debug)] struct MyError(String);
52///
53/// impl std::fmt::Display for MyError {
54/// fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { write!(f, "{}", self.0) }
55/// }
56/// impl std::error::Error for MyError {}
57/// impl From<String> for MyError { fn from(s: String) -> Self { MyError(s) } }
58///
59/// #[async_trait]
60/// impl ActorEntity for MyEntity {
61/// type Id = u32;
62/// type Create = MyCreate;
63/// type Update = MyUpdate;
64/// type Action = MyAction;
65/// type ActionResult = ();
66/// type Context = (); // No dependencies in this example
67/// type Error = MyError;
68///
69/// fn from_create_params(id: u32, _: MyCreate) -> Result<Self, Self::Error> { Ok(Self { id }) }
70/// async fn on_update(&mut self, _: MyUpdate, _: &()) -> Result<(), Self::Error> { Ok(()) }
71/// async fn handle_action(&mut self, _: MyAction, _: &()) -> Result<(), Self::Error> { Ok(()) }
72/// }
73///
74/// #[tokio::main]
75/// async fn main() {
76/// // 1. Create
77/// let (actor, client) = ResourceActor::<MyEntity>::new(10);
78///
79/// // 2. Wire & Run
80/// tokio::spawn(actor.run(()));
81///
82/// // 3. Use
83/// let _ = client.create(MyCreate).await;
84/// }
85/// ```
86///
87/// # Implementation Details
88///
89/// The actor maintains an internal `HashMap` (`store`) mapping IDs to entities and a `u32` counter (`next_id`) for ID generation.
90///
91/// ## Operations
92///
93/// * **Create**:
94/// 1. Generates a new ID using the internal `next_id` counter (incrementing it).
95/// 2. Converts the `u32` ID to `T::Id`.
96/// 3. Calls `T::from_create_params` to instantiate the entity.
97/// 4. Calls the `on_create` lifecycle hook.
98/// 5. Inserts the new entity into the `store`.
99/// 6. Returns the new ID.
100///
101/// * **Get**:
102/// 1. Looks up the entity in the `store` by ID.
103/// 2. Returns a clone of the entity if found, or `None`.
104///
105/// * **Update**:
106/// 1. Looks up the entity in the `store` (mutable access).
107/// 2. Calls the `on_update` lifecycle hook with the update DTO.
108/// 3. The entity modifies its own state within the hook.
109/// 4. Returns the updated entity state.
110///
111/// * **Delete**:
112/// 1. Looks up the entity in the `store`.
113/// 2. Calls the `on_delete` lifecycle hook.
114/// 3. Removes the entity from the `store`.
115///
116/// * **Action**:
117/// 1. Looks up the entity in the `store` (mutable access).
118/// 2. Calls the `handle_action` hook with the custom action enum.
119/// 3. Returns the result of the action.
120pub struct ResourceActor<T: ActorEntity> {
121 receiver: mpsc::Receiver<ResourceRequest<T>>,
122 store: HashMap<T::Id, T>,
123 next_id: u32,
124}
125
126impl<T: ActorEntity> ResourceActor<T> {
127 /// Creates a new `ResourceActor` and its associated `ResourceClient`.
128 ///
129 /// # Arguments
130 ///
131 /// * `buffer_size` - The capacity of the MPSC channel. If the channel is full,
132 /// calls to the client will wait until there is space.
133 ///
134 /// # Returns
135 ///
136 /// A tuple containing:
137 /// 1. The `ResourceActor` instance (the server), which must be run via `.run()`.
138 /// 2. The `ResourceClient` instance, which can be cloned and shared to send requests.
139 pub fn new(buffer_size: usize) -> (Self, ResourceClient<T>) {
140 let (sender, receiver) = mpsc::channel(buffer_size);
141 let actor = Self {
142 receiver,
143 store: HashMap::new(),
144 next_id: 1,
145 };
146 let client = ResourceClient::new(sender);
147 (actor, client)
148 }
149
150 /// Runs the actor's event loop, processing messages until the channel closes.
151 ///
152 /// # Context Injection
153 /// The `context` argument is injected into every entity hook. This allows entities
154 /// to access external dependencies (like other clients) that were created *after*
155 /// the actor was instantiated but *before* the loop started.
156 pub async fn run(mut self, context: T::Context) {
157 // Extract just the type name (e.g., "User" instead of "actor_recipe::model::user::User")
158 let entity_type = std::any::type_name::<T>()
159 .split("::")
160 .last()
161 .unwrap_or("Unknown");
162 info!(entity_type, "Actor started");
163
164 while let Some(msg) = self.receiver.recv().await {
165 match msg {
166 ResourceRequest::Create { params, respond_to } => {
167 debug!(entity_type, ?params, "Create");
168 let id = T::Id::from(self.next_id);
169 self.next_id += 1;
170
171 match T::from_create_params(id.clone(), params) {
172 Ok(mut item) => {
173 // Await the async hook
174 if let Err(e) = item.on_create(&context).await {
175 warn!(entity_type, error = %e, "on_create failed");
176 let _ =
177 respond_to.send(Err(FrameworkError::EntityError(Box::new(e))));
178 continue;
179 }
180 self.store.insert(id.clone(), item);
181 info!(entity_type, %id, size = self.store.len(), "Created");
182 let _ = respond_to.send(Ok(id));
183 }
184 Err(e) => {
185 warn!(entity_type, error = %e, "Create failed");
186 let _ = respond_to.send(Err(FrameworkError::EntityError(Box::new(e))));
187 }
188 }
189 }
190 ResourceRequest::Get { id, respond_to } => {
191 let item = self.store.get(&id).cloned();
192 let found = item.is_some();
193 debug!(entity_type, %id, found, "Get");
194 let _ = respond_to.send(Ok(item));
195 }
196 ResourceRequest::Update {
197 id,
198 update,
199 respond_to,
200 } => {
201 debug!(entity_type, %id, ?update, "Update");
202 if let Some(item) = self.store.get_mut(&id) {
203 // Await the async hook
204 if let Err(e) = item.on_update(update, &context).await {
205 warn!(entity_type, %id, error = %e, "Update failed");
206 let _ = respond_to.send(Err(FrameworkError::EntityError(Box::new(e))));
207 continue;
208 }
209 info!(entity_type, %id, "Updated");
210 let _ = respond_to.send(Ok(item.clone()));
211 } else {
212 warn!(entity_type, %id, "Not found");
213 let _ = respond_to.send(Err(FrameworkError::NotFound(id.to_string())));
214 }
215 }
216 ResourceRequest::Delete { id, respond_to } => {
217 debug!(entity_type, %id, "Delete");
218 if let Some(item) = self.store.get(&id) {
219 // Await the async hook
220 if let Err(e) = item.on_delete(&context).await {
221 warn!(entity_type, %id, error = %e, "on_delete failed");
222 let _ = respond_to.send(Err(FrameworkError::EntityError(Box::new(e))));
223 continue;
224 }
225 self.store.remove(&id);
226 info!(entity_type, %id, size = self.store.len(), "Deleted");
227 let _ = respond_to.send(Ok(()));
228 } else {
229 warn!(entity_type, %id, "Not found");
230 let _ = respond_to.send(Err(FrameworkError::NotFound(id.to_string())));
231 }
232 }
233 ResourceRequest::Action {
234 id,
235 action,
236 respond_to,
237 } => {
238 debug!(entity_type, %id, ?action, "Action");
239 if let Some(item) = self.store.get_mut(&id) {
240 // Await the async hook
241 let result = item
242 .handle_action(action, &context)
243 .await
244 .map_err(|e| FrameworkError::EntityError(Box::new(e)));
245 match &result {
246 Ok(_) => info!(entity_type, %id, "Action ok"),
247 Err(e) => warn!(entity_type, %id, error = %e, "Action failed"),
248 }
249 let _ = respond_to.send(result);
250 } else {
251 warn!(entity_type, %id, "Not found");
252 let _ = respond_to.send(Err(FrameworkError::NotFound(id.to_string())));
253 }
254 }
255 }
256 }
257
258 info!(entity_type, size = self.store.len(), "Shutdown");
259 }
260}