Skip to main content

renox_core/
events.rs

1//! Events and listeners, for decoupling modules: the order module emits
2//! `OrderPlaced`, and the stock and mail modules react to it.
3//!
4//! ```
5//! # use renox::prelude::*;
6//! # use serde::{Deserialize, Serialize};
7//! # #[derive(Serialize, Deserialize)] struct SendReceipt { order_id: i64 }
8//! # impl Job for SendReceipt { const NAME: &'static str = "send-receipt"; async fn handle(self, _: JobContext) -> Result { Ok(()) } }
9//! #[derive(Clone)]
10//! struct OrderPlaced { order_id: i64 }
11//! impl Event for OrderPlaced {}
12//!
13//! # let _ =
14//! App::new().listen(|event: OrderPlaced, state| async move {
15//!     state.dispatch(SendReceipt { order_id: event.order_id }).await?; // slow work: queue it
16//!     Ok(())
17//! })
18//! # ;
19//!
20//! # async fn demo(state: AppState, order_id: i64) -> Result {
21//! state.emit(OrderPlaced { order_id }).await?;
22//! # Ok(()) }
23//! ```
24
25use std::any::{Any, TypeId};
26use std::collections::HashMap;
27use std::future::Future;
28use std::pin::Pin;
29use std::sync::Arc;
30
31use crate::{AppState, Error, Result};
32
33/// Something that happened. Listeners get a clone each.
34pub trait Event: Clone + Send + Sync + 'static {}
35
36pub(crate) type ListenerFn = Arc<
37    dyn Fn(Box<dyn Any + Send>, AppState) -> Pin<Box<dyn Future<Output = Result> + Send>>
38        + Send
39        + Sync,
40>;
41
42pub(crate) type Listeners = Arc<HashMap<TypeId, Vec<ListenerFn>>>;
43
44pub(crate) fn listener<E, F, Fut>(listener: F) -> (TypeId, ListenerFn)
45where
46    E: Event,
47    F: Fn(E, AppState) -> Fut + Send + Sync + 'static,
48    Fut: Future<Output = Result> + Send + 'static,
49{
50    let run: ListenerFn = Arc::new(move |event, state| match event.downcast::<E>() {
51        Ok(event) => Box::pin(listener(*event, state)),
52        Err(_) => Box::pin(async { Ok(()) }),
53    });
54    (TypeId::of::<E>(), run)
55}
56
57impl AppState {
58    /// Runs every listener of `E` in registration order: those added with
59    /// `App::listen` first, then modules' (registered at boot). All of them
60    /// run even if one fails; the first error is returned.
61    pub async fn emit<E: Event>(&self, event: E) -> Result {
62        // `TestApp::fake_events`: record it, run nothing.
63        if self.fakes.record_event(event.clone()) {
64            return Ok(());
65        }
66        let Some(listeners) = self.listeners.get(&TypeId::of::<E>()) else {
67            return Ok(());
68        };
69        let mut first_error: Option<Error> = None;
70        for listener in listeners {
71            // A panicking listener is a failed one; the others still run.
72            let run = std::panic::AssertUnwindSafe(listener(Box::new(event.clone()), self.clone()));
73            let outcome = futures_util::FutureExt::catch_unwind(run)
74                .await
75                .unwrap_or_else(|_| Err(anyhow::anyhow!("the listener panicked").into()));
76            if let Err(err) = outcome {
77                tracing::error!(event = std::any::type_name::<E>(), error = ?err, "listener failed");
78                first_error.get_or_insert(err);
79            }
80        }
81        first_error.map_or(Ok(()), Err)
82    }
83}
84
85#[cfg(test)]
86mod tests {
87    use super::*;
88
89    #[derive(Clone)]
90    struct Placed;
91    impl Event for Placed {}
92
93    #[derive(Clone)]
94    struct Cancelled;
95    impl Event for Cancelled {}
96
97    /// Listeners are kept per event type, so one never gets another type;
98    /// if it did, it would do nothing rather than fail.
99    #[tokio::test]
100    async fn a_listener_ignores_an_event_of_another_type() {
101        let app = crate::testing::TestApp::new(crate::App::new()).await;
102        let (kind, run) = listener(|_: Placed, _state| async {
103            Err(Error::Internal(anyhow::anyhow!("only for Placed")))
104        });
105        assert_eq!(kind, TypeId::of::<Placed>());
106        assert!(run(Box::new(Cancelled), app.state().clone()).await.is_ok());
107        assert!(run(Box::new(Placed), app.state().clone()).await.is_err());
108    }
109}