Skip to main content

apalis_workflow/sequential/
step.rs

1use apalis_core::{
2    backend::{Backend, WireFormatBackend},
3    error::BoxDynError,
4};
5
6use crate::sequential::router::WorkflowRouter;
7
8/// A layer to wrap a step
9pub trait Layer<S> {
10    /// The resulting step type after layering.
11    type Step;
12    /// Wrap the given step with this layer.
13    fn layer(&self, step: S) -> Self::Step;
14}
15
16/// A sequential step
17///
18/// A single unit of work in a sequential workflow pipeline.
19pub trait Step<Input, B>
20where
21    B: Backend + WireFormatBackend,
22{
23    /// The response type produced by the step.
24    type Response;
25    /// The error type produced by the step.
26    type Error;
27
28    /// Register the step with the workflow router.
29    fn register(&mut self, router: &mut WorkflowRouter<B>) -> Result<(), BoxDynError>;
30}
31
32/// A no-op identity layer.
33#[derive(Clone, Debug)]
34pub struct Identity;
35
36impl<S> Layer<S> for Identity {
37    type Step = S;
38
39    fn layer(&self, step: S) -> Self::Step {
40        step
41    }
42}
43
44/// Two steps chained together.
45#[derive(Clone, Debug)]
46pub struct Stack<Inner, Outer> {
47    inner: Inner,
48    outer: Outer,
49}
50impl<Inner, Outer> Stack<Inner, Outer> {
51    /// Create a new `Stack`.
52    pub const fn new(inner: Inner, outer: Outer) -> Self {
53        Self { inner, outer }
54    }
55}
56
57impl<S, Inner, Outer> Layer<S> for Stack<Inner, Outer>
58where
59    Inner: Layer<S>,
60    Outer: Layer<Inner::Step>,
61{
62    type Step = Outer::Step;
63
64    fn layer(&self, service: S) -> Self::Step {
65        let inner = self.inner.layer(service);
66
67        self.outer.layer(inner)
68    }
69}