diff --git a/docs/superpowers/plans/2026-07-28-rust-async-course.md b/docs/superpowers/plans/2026-07-28-rust-async-course.md index 1108f2f..b39ea8a 100644 --- a/docs/superpowers/plans/2026-07-28-rust-async-course.md +++ b/docs/superpowers/plans/2026-07-28-rust-async-course.md @@ -84,7 +84,7 @@ Para cada capítulo, antes de pasar al siguiente: - [ ] Capítulo 10: modelo de actores. - [x] #34 Especificar propiedad, mensajes, fallas y alternativas. - - [ ] #36 Implementar y probar un modelo educativo mínimo. + - [x] #36 Implementar y probar un modelo educativo mínimo. - [ ] #38 Escribir capítulo, diagrama, ejemplos y ejercicios. - [ ] Completar ruta de lectura, glosario, referencias cruzadas y verificación final de coherencia del curso. diff --git a/src/actor_model.rs b/src/actor_model.rs new file mode 100644 index 0000000..05dacaa --- /dev/null +++ b/src/actor_model.rs @@ -0,0 +1,35 @@ +//! Modelo educativo de un actor que posee un contador. + +use tokio::sync::{mpsc, oneshot}; +use tokio::task::JoinHandle; + +/// Mensajes que entiende el actor contador. +pub enum CounterMessage { + /// Suma una cantidad al estado que posee el actor. + Increment(u64), + /// Devuelve el valor observado por el actor al procesar la consulta. + Get(oneshot::Sender), +} + +/// Inicia un actor contador y devuelve su buzón y la tarea que lo ejecuta. +/// +/// El actor termina al cerrar todos los emisores. `capacity` conserva el +/// backpressure de un canal acotado. +#[must_use] +pub fn spawn_counter(capacity: usize) -> (mpsc::Sender, JoinHandle<()>) { + let (sender, mut receiver) = mpsc::channel(capacity); + let task = tokio::spawn(async move { + let mut count = 0_u64; + + while let Some(message) = receiver.recv().await { + match message { + CounterMessage::Increment(amount) => count += amount, + CounterMessage::Get(reply) => { + let _ = reply.send(count); + } + } + } + }); + + (sender, task) +} diff --git a/src/lib.rs b/src/lib.rs index bb34ad2..4718b8f 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -6,6 +6,7 @@ #![forbid(unsafe_code)] +pub mod actor_model; pub mod async_channels; pub mod cooperative; pub mod coordination; diff --git a/tests/actor_model_test.rs b/tests/actor_model_test.rs new file mode 100644 index 0000000..59a5c88 --- /dev/null +++ b/tests/actor_model_test.rs @@ -0,0 +1,51 @@ +use rust_async::actor_model::{spawn_counter, CounterMessage}; +use tokio::sync::oneshot; + +#[tokio::test] +async fn actor_applies_messages_in_order_and_owns_its_state() { + let (sender, task) = spawn_counter(2); + sender + .send(CounterMessage::Increment(3)) + .await + .expect("actor should accept the message"); + sender + .send(CounterMessage::Increment(4)) + .await + .expect("actor should accept the message"); + + let (reply_sender, reply_receiver) = oneshot::channel(); + sender + .send(CounterMessage::Get(reply_sender)) + .await + .expect("actor should accept the query"); + + assert_eq!(reply_receiver.await.expect("actor should reply"), 7); + + drop(sender); + task.await.expect("actor task should finish cleanly"); +} + +#[tokio::test] +async fn actor_can_reply_to_more_than_one_query() { + let (sender, task) = spawn_counter(3); + let (first_reply_sender, first_reply_receiver) = oneshot::channel(); + sender + .send(CounterMessage::Get(first_reply_sender)) + .await + .expect("actor should accept the query"); + assert_eq!(first_reply_receiver.await.expect("actor should reply"), 0); + + sender + .send(CounterMessage::Increment(1)) + .await + .expect("actor should accept the message"); + let (second_reply_sender, second_reply_receiver) = oneshot::channel(); + sender + .send(CounterMessage::Get(second_reply_sender)) + .await + .expect("actor should accept the query"); + assert_eq!(second_reply_receiver.await.expect("actor should reply"), 1); + + drop(sender); + task.await.expect("actor task should finish cleanly"); +}