From 2fa62b0c07826dbf9f51725b4bc0479e06af4ad2 Mon Sep 17 00:00:00 2001 From: Joel Alvarez Date: Tue, 28 Jul 2026 13:13:24 -0700 Subject: [PATCH] feat: add async channel models --- .../plans/2026-07-28-rust-async-course.md | 2 ++ src/async_channels.rs | 36 +++++++++++++++++++ src/lib.rs | 1 + tests/async_channels_test.rs | 18 ++++++++++ 4 files changed, 57 insertions(+) create mode 100644 src/async_channels.rs create mode 100644 tests/async_channels_test.rs 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 a3fcdab..734599d 100644 --- a/docs/superpowers/plans/2026-07-28-rust-async-course.md +++ b/docs/superpowers/plans/2026-07-28-rust-async-course.md @@ -77,6 +77,8 @@ Para cada capítulo, antes de pasar al siguiente: - [x] #28 Escribir capítulo, diagrama, ejemplos y ejercicios. - [ ] Capítulo 09: canales y sincronización asíncrona. - [x] #29 Especificar backpressure, cierre y sincronización. + - [x] #30 Implementar y probar modelos de coordinación. + - [ ] #32 Escribir capítulo, diagrama, ejemplos y ejercicios. ### Milestone 4: Composición avanzada diff --git a/src/async_channels.rs b/src/async_channels.rs new file mode 100644 index 0000000..6dfc999 --- /dev/null +++ b/src/async_channels.rs @@ -0,0 +1,36 @@ +//! Modelos mínimos de canales acotados con Tokio. + +use tokio::sync::mpsc; + +/// Comprueba que un canal acotado comunica la falta de capacidad. +/// +/// La función usa `try_send` para observar presión de cola sin introducir una +/// espera temporal: el segundo mensaje no cabe mientras nadie reciba el primero. +pub fn bounded_channel_applies_backpressure() -> bool { + let (sender, _receiver) = mpsc::channel(1); + + sender.try_send(1).is_ok() && sender.try_send(2).is_err() +} + +/// Envía mensajes, cierra todos los emisores y drena la cola en orden FIFO. +pub async fn drain_after_senders_close() -> Vec { + let (sender, mut receiver) = mpsc::channel(2); + sender.send(1).await.expect("el receptor sigue abierto"); + sender.send(2).await.expect("el receptor sigue abierto"); + drop(sender); + + let mut received = Vec::new(); + while let Some(message) = receiver.recv().await { + received.push(message); + } + + received +} + +/// Indica que un emisor no puede confirmar entrega tras el cierre del receptor. +pub async fn receiver_closure_rejects_send() -> bool { + let (sender, receiver) = mpsc::channel(1); + drop(receiver); + + sender.send(1).await.is_err() +} diff --git a/src/lib.rs b/src/lib.rs index 924d738..bb34ad2 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -6,6 +6,7 @@ #![forbid(unsafe_code)] +pub mod async_channels; pub mod cooperative; pub mod coordination; pub mod educational_future; diff --git a/tests/async_channels_test.rs b/tests/async_channels_test.rs new file mode 100644 index 0000000..cca0c3e --- /dev/null +++ b/tests/async_channels_test.rs @@ -0,0 +1,18 @@ +use rust_async::async_channels::{ + bounded_channel_applies_backpressure, drain_after_senders_close, receiver_closure_rejects_send, +}; + +#[tokio::test] +async fn bounded_channel_reports_full_queue() { + assert!(bounded_channel_applies_backpressure()); +} + +#[tokio::test] +async fn receiver_drains_messages_before_observing_sender_closure() { + assert_eq!(drain_after_senders_close().await, vec![1, 2]); +} + +#[tokio::test] +async fn closed_receiver_rejects_new_messages() { + assert!(receiver_closure_rejects_send().await); +}