Skip to main content

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}