hexeract_outbox/
handler.rs1use hexeract_core::HandlerContext;
2
3use crate::Event;
4use crate::OutboxError;
5
6#[trait_variant::make(Send)]
52pub trait Handler<E: Event>: Send + Sync + 'static {
53 type Error: Into<OutboxError> + Send + Sync + 'static;
55
56 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}