Skip to main content

apalis_core/backend/ext/
mod.rs

1//! Extension traits and combinators for [`Backend`] implementations.
2//!
3//! This module provides additional functionality for [`Backend`] implementations,
4//! including:
5//!
6//! - [`BackendExt`]: Extension methods for transforming, composing, and managing
7//!   backends.
8//! - [`InspectErr`]: A wrapper that allows inspection of errors produced by a backend.
9//! - [`MapErr`]: A wrapper that maps backend errors from one type to another.
10//! - [`Pipe`]: A utility for piping tasks from one backend to another.
11//! - [`BeforeStart`]: A lifecycle wrapper that runs an action before the backend starts.
12//! - [`AfterStart`]: A lifecycle wrapper that runs an action after the backend starts.
13//! - [`BeforeStop`]: A lifecycle wrapper that runs an action before the backend stops.
14//! - [`AfterStop`]: A lifecycle wrapper that runs an action after the backend stops.
15use std::{
16    task::{Context, Poll},
17    time::Duration,
18};
19
20use futures_core::Stream;
21
22#[cfg(feature = "tracing")]
23use crate::backend::ext::instrument::Instrumented;
24use crate::{
25    backend::{
26        Backend, BackendConfig, WireFormatBackend,
27        codec::Codec,
28        ext::{
29            inspect_err::InspectErr,
30            interleave::Interleave,
31            lifecycle::{AfterStart, AfterStop, BeforeStart, BeforeStop},
32            map_err::MapErr,
33            pipe::Pipe,
34            poll_strategy::{PollStrategy, PollWith, StreamStrategy},
35            wake_on_push::WakeOnPush,
36            with_codec::WithCodec,
37        },
38    },
39    error::BoxDynError,
40    task::Task,
41    worker::context::WorkerContext,
42};
43
44#[cfg(feature = "shared")]
45use crate::backend::ext::shared::Shared;
46
47#[cfg(feature = "sleep")]
48use crate::backend::ext::poll_strategy::{BackoffConfig, BackoffStrategy, IntervalStrategy};
49
50#[macro_use]
51pub mod delegate;
52/// A wrapper that allows inspecting errors produced by a backend.
53pub mod inspect_err;
54
55/// Extension allowing backends to be instrumented with a [tracing::Span].
56#[cfg(feature = "tracing")]
57pub mod instrument;
58/// A wrapper that allows merging a backend with a stream
59pub mod interleave;
60
61/// A wrapper that allows mapping the error type of a backend from `Self::Error` to another error type `E2`.
62pub mod map_err;
63pub mod pipe;
64pub mod poll_strategy;
65/// A wrapper that wakes the worker when a new item is fetched.
66pub mod wake_on_push;
67pub mod with_codec;
68
69pub mod lifecycle;
70
71/// A wrapper that makes a backend clonable
72#[cfg(feature = "shared")]
73pub mod shared;
74
75/// A wrapper that allows a backend to be used as a stream of tasks, without needing to know the concrete backend type at compile time.
76#[derive(Debug, thiserror::Error)]
77#[non_exhaustive]
78pub enum PollNextArgsError<B: Backend> {
79    /// The backend produced an error while polling for the next task.
80    #[error("backend error: {0}")]
81    BackendError(B::Error),
82    /// The backend produced a task, but the task's arguments could not be decoded.
83    #[error("failed to decode task args: {0}")]
84    DecodeError(BoxDynError),
85}
86
87/// Extension trait for `Backend` that provides additional combinators and utilities.
88pub trait BackendExt: Backend {
89    /// A convenience method for calling `poll_next` and decoding the `Args` in one step,
90    /// returning a `Task<Self::Args, ..>` instead of `Task<Self::Compact, ..>`.
91    #[allow(clippy::type_complexity)]
92    fn poll_next_args(
93        &mut self,
94        cx: &mut Context<'_>,
95        worker: &WorkerContext,
96    ) -> Poll<Option<Result<Task<Self::Args>, PollNextArgsError<Self>>>>
97    where
98        Self: Sized + BackendConfig + WireFormatBackend + Backend<Task = Task<Self::Compact>>,
99        Self::Codec: Codec<Self::Args, Compact = Self::Compact>,
100        <Self::Codec as Codec<Self::Args>>::Error: std::error::Error + Send + Sync + 'static,
101    {
102        let next = self.poll_next(cx, worker);
103        let codec = self.codec();
104        next.map(move |item| match item {
105            Some(Ok(task)) => {
106                let task = task.try_map_args(|compact| codec.decode(&compact));
107                Some(task.map_err(|e| PollNextArgsError::DecodeError(e.into())))
108            }
109            Some(Err(e)) => Some(Err(PollNextArgsError::BackendError(e))),
110            None => None,
111        })
112    }
113
114    /// Pipes every task polled from this backend into `sink`
115    ///
116    /// Useful for bridging two backend implementations — e.g. draining an
117    /// ephemeral/legacy queue into a durable one, or fanning a lightweight
118    /// source into a shared sink that multiple producers write into.
119    fn pipe_to<Dst>(self, backend: Dst) -> Pipe<Dst, Self>
120    where
121        Self: Sized,
122    {
123        Pipe::new(self, backend)
124    }
125
126    /// Attaches a callback `F` to be run on each error produced while polling the backend.
127    fn inspect_err<F>(self, f: F) -> InspectErr<Self, F>
128    where
129        Self: Sized,
130        F: Fn(&Self::Error),
131    {
132        InspectErr { backend: self, f }
133    }
134
135    /// Maps errors produced by the backend from `Self::Error` into `E2`, useful for
136    /// heterogeneous composed backends.
137    fn map_err<F, E2>(self, f: F) -> MapErr<Self, F>
138    where
139        Self: Sized,
140        F: Fn(Self::Error) -> E2,
141    {
142        MapErr { backend: self, f }
143    }
144
145    /// Swaps out the backend's serialization codec entirely (JSON,
146    /// MessagePack, Protobuf, ...) without touching storage logic.
147    fn with_codec<NewCodec>(self, codec: NewCodec) -> WithCodec<Self, NewCodec>
148    where
149        Self: Sized + BackendConfig,
150        NewCodec: Codec<Self::Args>,
151    {
152        WithCodec::new(self, codec)
153    }
154
155    /// Wake the worker when a stream receives a new item
156    fn poll_with_stream<S>(self, stream: S) -> PollWith<Self, StreamStrategy<S>>
157    where
158        Self: Sized,
159        S: Stream + Unpin + Send + 'static,
160    {
161        let strategy = StreamStrategy::new(stream);
162        PollWith::new(self, strategy)
163    }
164
165    /// Wake the worker periodically
166    #[cfg(feature = "sleep")]
167    fn poll_with_interval(self, duration: Duration) -> PollWith<Self, IntervalStrategy>
168    where
169        Self: Sized,
170    {
171        let strategy = IntervalStrategy::new(duration);
172        PollWith::new(self, strategy)
173    }
174
175    /// Wake the worker periodically with a backoff
176    #[cfg(feature = "sleep")]
177    fn poll_with_backoff(
178        self,
179        interval: Duration,
180        config: BackoffConfig,
181    ) -> PollWith<Self, BackoffStrategy>
182    where
183        Self: Sized,
184    {
185        let strategy = IntervalStrategy::new(interval).with_backoff(config);
186        PollWith::new(self, strategy)
187    }
188
189    /// Wake the worker with a custom strategy
190    fn poll_with_strategy<S>(self, strategy: S) -> PollWith<Self, S>
191    where
192        Self: Sized,
193        S: PollStrategy,
194    {
195        PollWith::new(self, strategy)
196    }
197
198    #[cfg(feature = "tracing")]
199    /// Provides a span to decorate emitted events
200    fn instrumented(self, span: tracing::Span) -> Instrumented<Self>
201    where
202        Self: Sized,
203    {
204        Instrumented::new(self, span)
205    }
206
207    /// Runs an async callback once, before the backend's first `poll_ready` is delegated.
208    fn before_start<F, Fut>(self, f: F) -> BeforeStart<Self, Self::Error>
209    where
210        Self: Sized,
211        F: Fn(&mut Self) -> Fut + Send + Sync + 'static,
212        Fut: Future<Output = Result<(), Self::Error>> + Send + 'static,
213    {
214        BeforeStart::new(self, f)
215    }
216
217    /// Runs an async callback once, before the backend's poll_close is called.
218    fn before_stop<F, Fut>(self, f: F) -> BeforeStop<Self, Self::Error>
219    where
220        Self: Sized,
221        F: Fn(&mut Self) -> Fut + Send + Sync + 'static,
222        Fut: Future<Output = Result<(), Self::Error>> + Send + 'static,
223    {
224        BeforeStop::new(self, f)
225    }
226
227    /// Runs an async callback once, after the backend's first `poll_ready` is successful.
228    fn after_start<F, Fut>(self, f: F) -> AfterStart<Self, Self::Error>
229    where
230        Self: Sized,
231        for<'c> F: Fn(&mut Self) -> Fut + Send + Sync + 'static,
232        Fut: Future<Output = Result<(), Self::Error>> + Send + 'static,
233    {
234        AfterStart::new(self, f)
235    }
236
237    /// Runs an async callback once, after the worker has stopped and backend has cleaned up.
238    fn after_stop<F, Fut>(self, f: F) -> AfterStop<Self, Self::Error>
239    where
240        Self: Sized,
241        F: Fn(&mut Self) -> Fut + Send + Sync + 'static,
242        Fut: Future<Output = Result<(), Self::Error>> + Send + 'static,
243    {
244        AfterStop::new(self, f)
245    }
246
247    /// Interleaves the external stream with the backend.
248    fn interleave<S>(self, stream: S) -> Interleave<Self, S>
249    where
250        Self: Sized,
251        S: Stream<Item = Result<Self::Task, Self::Error>> + Unpin,
252    {
253        Interleave::new(self, stream)
254    }
255
256    /// Wakes the worker when a new item is pushed
257    fn wake_on_push(self) -> WakeOnPush<Self>
258    where
259        Self: Sized,
260    {
261        WakeOnPush::new(self)
262    }
263
264    /// Create a cloneable handle to the inner backend where all handles are clone.
265    #[cfg(feature = "shared")]
266    fn shared(self) -> Shared<Self>
267    where
268        Self: WireFormatBackend + Send,
269        Self::Codec: Clone,
270    {
271        Shared::new(self)
272    }
273}
274
275impl<B: Backend> BackendExt for B {}