apalis_core/worker/ext/ack/mod.rs
1//! Traits and utilities for acknowledging task completion
2//!
3//! The [`Acknowledge`] trait and related types are responsible for adding custom
4//! acknowledgment logic to workers. You can use [`AcknowledgeLayer`] to wrap
5//! a worker service and invoke your acknowledgment handler after each task execution.
6//!
7//! # Example
8//!
9//! ```rust
10//! # use apalis_core::worker::{builder::WorkerBuilder, ext::ack::{Acknowledge, AcknowledgeLayer}};
11//! # use apalis_core::backend::memory::MemoryStorage;
12//! # use apalis_core::worker::context::WorkerContext;
13//! # use apalis_core::task::ExecutionContext;
14//! # use apalis_core::error::BoxDynError;
15//! # use futures_util::{future::{ready, BoxFuture}, FutureExt};
16//! # use std::fmt::Debug;
17//! # use tokio::sync::mpsc::error::SendError;
18//! # use apalis_core::worker::ext::ack::AcknowledgementExt;
19//! # use apalis_core::backend::TaskSink;
20//! # use crate::apalis_core::worker::ext::event_listener::EventListenerExt;
21//!
22//! #[tokio::main]
23//! async fn main() {
24//! let mut in_memory = MemoryStorage::new();
25//! in_memory.push(42).await.unwrap();
26//!
27//! async fn task(
28//! task: u32,
29//! worker: WorkerContext,
30//! ) -> Result<(), BoxDynError> {
31//! # worker.stop().unwrap();
32//! Ok(())
33//! }
34//!
35//! #[derive(Debug, Clone)]
36//! struct MyAcknowledger;
37//!
38//! impl Acknowledge<()> for MyAcknowledger {
39//! type Error = SendError<()>;
40//! type Future = BoxFuture<'static, Result<(), Self::Error>>;
41//! fn ack(
42//! &mut self,
43//! res: &Result<(), BoxDynError>,
44//! ctx: &ExecutionContext,
45//! ) -> Self::Future {
46//! println!("{res:?}, {ctx:?}");
47//! ready(Ok(())).boxed()
48//! }
49//! }
50//!
51//! let worker = WorkerBuilder::new("rango-tango")
52//! .backend(in_memory)
53//! .ack_with(MyAcknowledger)
54//! .on_event(|worker, ev| {
55//! println!("On Event = {:?}", ev);
56//! })
57//! .build(task);
58//! worker.run().await.unwrap();
59//! }
60//! ```
61use futures_util::future::BoxFuture;
62use std::{future::Future, task::Poll};
63use tower_layer::{Layer, Stack};
64use tower_service::Service;
65
66use crate::{
67 backend::Backend,
68 error::BoxDynError,
69 task::{ExecutionContext, Task},
70 worker::builder::WorkerBuilder,
71};
72
73/// Extension trait for adding acknowledgment handling to workers
74///
75/// See [module level documentation](self) for more details.
76pub trait AcknowledgementExt<Args, Source, Middleware, Ack, Res>: Sized
77where
78 Source: Backend,
79 Ack: Acknowledge<Res>,
80{
81 /// Add an acknowledgment handler to the worker
82 fn ack_with(
83 self,
84 ack: Ack,
85 ) -> WorkerBuilder<Args, Source, Stack<AcknowledgeLayer<Ack>, Middleware>>;
86}
87
88/// Acknowledge the result of a task processing
89///
90/// See [module level documentation](self) for more details.
91pub trait Acknowledge<Res> {
92 /// The error type returned by the acknowledgment process
93 type Error;
94 /// The future returned by the `ack` method
95 type Future: Future<Output = Result<(), Self::Error>>;
96 /// Acknowledge the result of a task processing
97 fn ack(&mut self, res: &Result<Res, BoxDynError>, ctx: &ExecutionContext) -> Self::Future;
98}
99
100impl<Res, F, Fut, E> Acknowledge<Res> for F
101where
102 for<'c> F: FnMut(&'c Result<Res, BoxDynError>, &'c ExecutionContext) -> Fut,
103 Fut: Future<Output = Result<(), E>>,
104{
105 type Error = E;
106 type Future = Fut;
107
108 fn ack(&mut self, res: &Result<Res, BoxDynError>, ctx: &ExecutionContext) -> Self::Future {
109 (self)(res, ctx)
110 }
111}
112
113/// Layer that adds acknowledgment functionality to services
114///
115/// See [module level documentation](self) for more details.
116#[derive(Debug, Clone)]
117pub struct AcknowledgeLayer<A> {
118 acknowledger: A,
119}
120
121impl<A> AcknowledgeLayer<A> {
122 /// Create a new acknowledgment layer
123 pub fn new(acknowledger: A) -> Self {
124 Self { acknowledger }
125 }
126}
127
128impl<S, A> Layer<S> for AcknowledgeLayer<A>
129where
130 A: Clone,
131{
132 type Service = AcknowledgeService<S, A>;
133
134 fn layer(&self, inner: S) -> Self::Service {
135 AcknowledgeService {
136 inner,
137 acknowledger: self.acknowledger.clone(),
138 }
139 }
140}
141
142/// Service that wraps another service and acknowledges task completion
143///
144/// See [module level documentation](self) for more details.
145
146#[derive(Debug, Clone)]
147pub struct AcknowledgeService<S, A> {
148 inner: S,
149 acknowledger: A,
150}
151
152impl<S, A, Args, Res> Service<Task<Args>> for AcknowledgeService<S, A>
153where
154 S: Service<Task<Args>, Response = Res>,
155 A: Acknowledge<Res> + Clone + Send + 'static,
156 S::Error: Into<BoxDynError>,
157 A::Error: std::error::Error + Send + Sync + 'static,
158 S::Future: Send + 'static,
159 A::Future: Send + 'static,
160 Res: Send,
161{
162 type Response = Res;
163 type Error = BoxDynError;
164 type Future = BoxFuture<'static, Result<Res, BoxDynError>>;
165
166 fn poll_ready(&mut self, cx: &mut std::task::Context<'_>) -> Poll<Result<(), Self::Error>> {
167 self.inner.poll_ready(cx).map_err(|e| e.into())
168 }
169
170 fn call(&mut self, task: Task<Args>) -> Self::Future {
171 let ctx = task.ctx().clone();
172 let future = self.inner.call(task);
173 let mut acknowledger = self.acknowledger.clone();
174 Box::pin(async move {
175 let res = future.await.map_err(|e| e.into());
176 acknowledger.ack(&res, &ctx).await?;
177 res
178 })
179 }
180}
181
182impl<Args, B, M, Ack, Res> AcknowledgementExt<Args, B, M, Ack, Res> for WorkerBuilder<Args, B, M>
183where
184 M: Layer<AcknowledgeLayer<Ack>>,
185 Ack: Acknowledge<Res>,
186 B: Backend,
187{
188 fn ack_with(self, ack: Ack) -> WorkerBuilder<Args, B, Stack<AcknowledgeLayer<Ack>, M>> {
189 let this = self.layer(AcknowledgeLayer::new(ack));
190 WorkerBuilder {
191 context: this.context,
192 request: this.request,
193 layer: this.layer,
194 source: this.source,
195 shutdown: this.shutdown,
196 event_handler: this.event_handler,
197 }
198 }
199}