Skip to main content

hexeract_outbox/
handler.rs

1use hexeract_core::HandlerContext;
2
3use crate::Event;
4use crate::OutboxError;
5
6/// Asynchronous handler dispatched by the outbox worker for each event of type `E`.
7///
8/// Implementors describe the side effect to perform when an event lands:
9/// write to an audit log, publish to a broker, send a notification, etc.
10/// The handler does not return a business value: success is the side
11/// effect itself.
12///
13/// # Idempotency
14///
15/// The outbox guarantees at-least-once delivery. Handlers MUST therefore
16/// be idempotent: the same event can be delivered more than once if a
17/// previous attempt crashed between the side effect and the database
18/// commit that marks the row as delivered.
19///
20/// # Example
21///
22/// ```
23/// use hexeract_core::HandlerContext;
24/// use hexeract_outbox::{Event, Handler, OutboxError};
25/// use serde::{Deserialize, Serialize};
26///
27/// #[derive(Debug, Serialize, Deserialize)]
28/// struct UserRegistered {
29///     user_id: uuid::Uuid,
30/// }
31///
32/// impl Event for UserRegistered {
33///     const EVENT_TYPE: &'static str = "users.registered";
34/// }
35///
36/// struct AuditWriter;
37///
38/// impl Handler<UserRegistered> for AuditWriter {
39///     type Error = OutboxError;
40///
41///     async fn handle(
42///         &self,
43///         event: UserRegistered,
44///         _ctx: &HandlerContext,
45///     ) -> Result<(), Self::Error> {
46///         let _ = event.user_id;
47///         Ok(())
48///     }
49/// }
50/// ```
51#[trait_variant::make(Send)]
52pub trait Handler<E: Event>: Send + Sync + 'static {
53    /// Handler-defined error type, convertible into [`OutboxError`].
54    type Error: Into<OutboxError> + Send + Sync + 'static;
55
56    /// Process the event and produce its side effect.
57    async fn handle(&self, event: E, ctx: &HandlerContext) -> Result<(), Self::Error>;
58}
59
60#[cfg(test)]
61mod tests {
62    use super::*;
63    use hexeract_core::CorrelationId;
64    use hexeract_core::MessageId;
65    use serde::Deserialize;
66    use serde::Serialize;
67    use std::sync::Arc;
68    use std::time::Duration;
69    use uuid::Uuid;
70
71    fn fresh_ctx() -> HandlerContext {
72        HandlerContext::new(MessageId::new(), CorrelationId::new())
73    }
74
75    fn assert_send<T: Send>(_: &T) {}
76
77    #[derive(Debug, Serialize, Deserialize)]
78    struct UserRegistered {
79        user_id: Uuid,
80    }
81
82    impl Event for UserRegistered {
83        const EVENT_TYPE: &'static str = "users.registered";
84    }
85
86    #[derive(Debug, Serialize, Deserialize)]
87    struct OrderPlaced {
88        order_id: Uuid,
89    }
90
91    impl Event for OrderPlaced {
92        const EVENT_TYPE: &'static str = "orders.placed";
93    }
94
95    #[derive(Debug, thiserror::Error)]
96    enum AuditError {
97        #[error("audit store unavailable")]
98        Unavailable,
99    }
100
101    impl From<AuditError> for OutboxError {
102        fn from(value: AuditError) -> Self {
103            Self::Internal(value.to_string())
104        }
105    }
106
107    struct AuditWriter {
108        accept: bool,
109    }
110
111    impl Handler<UserRegistered> for AuditWriter {
112        type Error = AuditError;
113        async fn handle(
114            &self,
115            _event: UserRegistered,
116            _ctx: &HandlerContext,
117        ) -> Result<(), Self::Error> {
118            if self.accept {
119                Ok(())
120            } else {
121                Err(AuditError::Unavailable)
122            }
123        }
124    }
125
126    #[tokio::test]
127    async fn handler_returns_unit_on_success() {
128        let writer = AuditWriter { accept: true };
129        let ctx = fresh_ctx();
130        writer
131            .handle(
132                UserRegistered {
133                    user_id: Uuid::nil(),
134                },
135                &ctx,
136            )
137            .await
138            .expect("handler should succeed");
139    }
140
141    #[tokio::test]
142    async fn handler_typed_error_converts_into_outbox_error() {
143        let writer = AuditWriter { accept: false };
144        let ctx = fresh_ctx();
145        let err = writer
146            .handle(
147                UserRegistered {
148                    user_id: Uuid::nil(),
149                },
150                &ctx,
151            )
152            .await
153            .expect_err("handler should fail");
154        assert!(matches!(err, AuditError::Unavailable));
155        let outbox_err: OutboxError = err.into();
156        assert!(matches!(outbox_err, OutboxError::Internal(_)));
157    }
158
159    struct DirectHandler;
160    impl Handler<UserRegistered> for DirectHandler {
161        type Error = OutboxError;
162        async fn handle(
163            &self,
164            _event: UserRegistered,
165            _ctx: &HandlerContext,
166        ) -> Result<(), Self::Error> {
167            Err(OutboxError::Internal("forced".into()))
168        }
169    }
170
171    #[tokio::test]
172    async fn handler_can_use_outbox_error_directly() {
173        let handler = DirectHandler;
174        let ctx = fresh_ctx();
175        let err = handler
176            .handle(
177                UserRegistered {
178                    user_id: Uuid::nil(),
179                },
180                &ctx,
181            )
182            .await
183            .expect_err("must fail");
184        assert!(matches!(err, OutboxError::Internal(_)));
185    }
186
187    #[tokio::test]
188    async fn handler_future_is_send() {
189        let writer = AuditWriter { accept: true };
190        let ctx = fresh_ctx();
191        let future = writer.handle(
192            UserRegistered {
193                user_id: Uuid::nil(),
194            },
195            &ctx,
196        );
197        assert_send(&future);
198        let _ = future.await;
199    }
200
201    #[tokio::test]
202    async fn handler_runs_in_spawned_task() {
203        let writer = Arc::new(AuditWriter { accept: true });
204        let cloned = Arc::clone(&writer);
205        let result = tokio::spawn(async move {
206            let ctx = fresh_ctx();
207            cloned
208                .handle(
209                    UserRegistered {
210                        user_id: Uuid::nil(),
211                    },
212                    &ctx,
213                )
214                .await
215        })
216        .await
217        .expect("task panicked");
218        assert!(result.is_ok());
219    }
220
221    struct EchoCtxHandler;
222    impl Handler<UserRegistered> for EchoCtxHandler {
223        type Error = OutboxError;
224        async fn handle(
225            &self,
226            _event: UserRegistered,
227            ctx: &HandlerContext,
228        ) -> Result<(), Self::Error> {
229            let _ = ctx.message_id;
230            let _ = ctx.correlation_id;
231            Ok(())
232        }
233    }
234
235    #[tokio::test]
236    async fn handler_reads_ids_from_context() {
237        let message_id = MessageId::new();
238        let correlation_id = CorrelationId::new();
239        let ctx = HandlerContext::new(message_id, correlation_id);
240
241        let handler = EchoCtxHandler;
242        handler
243            .handle(
244                UserRegistered {
245                    user_id: Uuid::nil(),
246                },
247                &ctx,
248            )
249            .await
250            .expect("handler should succeed");
251
252        assert_eq!(ctx.message_id, message_id);
253        assert_eq!(ctx.correlation_id, correlation_id);
254    }
255
256    struct SleepHandler;
257    impl Handler<UserRegistered> for SleepHandler {
258        type Error = OutboxError;
259        async fn handle(
260            &self,
261            _event: UserRegistered,
262            ctx: &HandlerContext,
263        ) -> Result<(), Self::Error> {
264            tokio::select! {
265                () = ctx.cancellation.cancelled() => Err(OutboxError::Internal("cancelled".into())),
266                () = tokio::time::sleep(Duration::from_millis(5_000)) => Ok(()),
267            }
268        }
269    }
270
271    #[tokio::test]
272    async fn handler_observes_external_cancellation() {
273        let ctx = fresh_ctx();
274        let token = ctx.cancellation.clone();
275
276        let handle = tokio::spawn(async move {
277            let handler = SleepHandler;
278            handler
279                .handle(
280                    UserRegistered {
281                        user_id: Uuid::nil(),
282                    },
283                    &ctx,
284                )
285                .await
286        });
287
288        tokio::time::sleep(Duration::from_millis(50)).await;
289        token.cancel();
290
291        let result = handle.await.expect("task panicked");
292        assert!(matches!(result, Err(OutboxError::Internal(ref m)) if m == "cancelled"));
293    }
294
295    #[tokio::test]
296    async fn handler_is_shareable_via_arc() {
297        let handler: Arc<AuditWriter> = Arc::new(AuditWriter { accept: true });
298        let h1 = Arc::clone(&handler);
299        let h2 = Arc::clone(&handler);
300
301        let t1 = tokio::spawn(async move {
302            let ctx = fresh_ctx();
303            h1.handle(
304                UserRegistered {
305                    user_id: Uuid::nil(),
306                },
307                &ctx,
308            )
309            .await
310        });
311        let t2 = tokio::spawn(async move {
312            let ctx = fresh_ctx();
313            h2.handle(
314                UserRegistered {
315                    user_id: Uuid::nil(),
316                },
317                &ctx,
318            )
319            .await
320        });
321
322        let (r1, r2) = tokio::join!(t1, t2);
323        assert!(r1.unwrap().is_ok());
324        assert!(r2.unwrap().is_ok());
325    }
326
327    struct MultiHandler;
328    impl Handler<UserRegistered> for MultiHandler {
329        type Error = OutboxError;
330        async fn handle(
331            &self,
332            _event: UserRegistered,
333            _ctx: &HandlerContext,
334        ) -> Result<(), Self::Error> {
335            Ok(())
336        }
337    }
338    impl Handler<OrderPlaced> for MultiHandler {
339        type Error = OutboxError;
340        async fn handle(
341            &self,
342            _event: OrderPlaced,
343            _ctx: &HandlerContext,
344        ) -> Result<(), Self::Error> {
345            Ok(())
346        }
347    }
348
349    #[tokio::test]
350    async fn one_struct_can_handle_multiple_event_types() {
351        let handler = MultiHandler;
352        let ctx = fresh_ctx();
353        Handler::<UserRegistered>::handle(
354            &handler,
355            UserRegistered {
356                user_id: Uuid::nil(),
357            },
358            &ctx,
359        )
360        .await
361        .expect("user handler must succeed");
362        let ctx = fresh_ctx();
363        Handler::<OrderPlaced>::handle(
364            &handler,
365            OrderPlaced {
366                order_id: Uuid::nil(),
367            },
368            &ctx,
369        )
370        .await
371        .expect("order handler must succeed");
372    }
373}