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}