Skip to main content

ntex_service/
pipeline.rs

1use std::{fmt, future, pin::Pin, task::Context, task::Poll};
2
3use crate::pl_inner::PipelineApi;
4use crate::{IntoService, Service, ServiceCaller, util::BoxFuture};
5
6pub use crate::pl_factory::PipelineFactory;
7pub use crate::pl_state::{PipelineState, PipelineStateBinding};
8
9/// Execution container for a service and its state.
10///
11/// A pipeline coordinates readiness, calls, and shutdown for the enclosed
12/// service chain.
13pub struct Pipeline<Req, Res, Err> {
14    api: PipelineApi<Req, Res, Err>,
15}
16
17/// An independently registered handle to a [`Pipeline`].
18///
19/// Bindings share the pipeline and its readiness state. Cloning a binding
20/// registers another handle.
21pub struct PipelineBinding<Req, Res, Err> {
22    idx: u32,
23    api: PipelineApi<Req, Res, Err>,
24}
25
26impl<Req, Res, Err> Pipeline<Req, Res, Err>
27where
28    Req: 'static,
29    Res: 'static,
30    Err: 'static,
31{
32    #[inline]
33    /// Creates a pipeline containing `service` and `state`.
34    pub fn new<S, St>(st: St, service: impl IntoService<S, St, Req>) -> Self
35    where
36        S: Service<St, Req, Res = Res, Error = Err> + 'static,
37        St: 'static,
38    {
39        Pipeline {
40            api: PipelineApi::new(service.into_service(), st),
41        }
42    }
43
44    #[inline]
45    /// Returns when the pipeline is ready to process requests.
46    ///
47    /// A successful check is consumed by the next call, which then skips
48    /// its own readiness check.
49    pub async fn ready(&self) -> Result<(), Err> {
50        future::poll_fn(|cx| self.api.poll_ready(cx)).await
51    }
52
53    #[inline]
54    /// Waits for readiness, then calls the service.
55    ///
56    /// The readiness check is skipped if the last pipeline readiness check
57    /// succeeded and no call has started since.
58    pub async fn call(&self, req: Req) -> Result<Res, Err> {
59        let pl = self.bind();
60        pl.api.call(pl.idx, req).await
61    }
62
63    #[inline]
64    /// Returns an owned future that waits for readiness and calls the service.
65    ///
66    /// Unlike [`Pipeline::call`], the returned future does not borrow the
67    /// pipeline and can be moved between local tasks. The readiness check is
68    /// skipped if the last pipeline readiness check succeeded and no call has
69    /// started since; this is decided when the future is first polled.
70    pub fn call_static(&self, req: Req) -> PipelineCall<Req, Res, Err> {
71        PipelineCall::new(self.bind(), req)
72    }
73
74    #[inline]
75    /// Returns `Ready` when the pipeline is ready to process requests.
76    ///
77    /// A successful check is consumed by the next call, which then skips
78    /// its own readiness check.
79    pub fn poll_ready(&self, cx: &mut Context<'_>) -> Poll<Result<(), Err>> {
80        self.api.poll_ready(cx)
81    }
82
83    #[inline]
84    /// Returns `Ready` when the service has been properly shut down.
85    pub fn poll_shutdown(&self, cx: &mut Context<'_>) -> Poll<()> {
86        self.api.poll_shutdown(cx)
87    }
88
89    #[inline]
90    /// Checks whether pipeline shutdown has been initiated.
91    pub fn is_shutdown(&self) -> bool {
92        self.api.is_shutdown()
93    }
94
95    #[inline]
96    /// Shuts down the enclosed service.
97    pub async fn shutdown(&self) {
98        future::poll_fn(|cx| self.api.poll_shutdown(cx)).await;
99    }
100
101    #[inline]
102    /// Creates a new binding to this pipeline.
103    ///
104    /// The binding can be used to check readiness and call the service.
105    pub fn bind(&self) -> PipelineBinding<Req, Res, Err> {
106        PipelineBinding::new(self)
107    }
108}
109
110impl<Req, Res, Err> ServiceCaller<Req, Res, Err> for Pipeline<Req, Res, Err>
111where
112    Req: 'static,
113    Res: 'static,
114    Err: 'static,
115{
116    #[inline]
117    async fn call_service(&self, req: Req) -> Result<Res, Err> {
118        let pl = self.bind();
119        pl.api.call(pl.idx, req).await
120    }
121}
122
123impl<Req, Res, Err> Drop for Pipeline<Req, Res, Err> {
124    #[inline]
125    fn drop(&mut self) {
126        self.api.unreg(0);
127    }
128}
129
130impl<Req, Res, Err> PipelineBinding<Req, Res, Err>
131where
132    Req: 'static,
133    Res: 'static,
134    Err: 'static,
135{
136    fn new(pl: &Pipeline<Req, Res, Err>) -> Self {
137        Self {
138            idx: pl.api.reg(),
139            api: pl.api.clone(),
140        }
141    }
142
143    pub(crate) fn with(idx: u32, api: PipelineApi<Req, Res, Err>) -> Self {
144        Self { idx, api }
145    }
146
147    #[inline]
148    /// Waits until the pipeline is ready to process a request.
149    ///
150    /// A successful check is consumed by the next call, which then skips
151    /// its own readiness check.
152    pub async fn ready(&self) -> Result<(), Err> {
153        self.api.ready(self.idx).await
154    }
155
156    #[inline]
157    /// Waits for readiness, then calls the service.
158    ///
159    /// The readiness check is skipped if the last pipeline readiness check
160    /// succeeded and no call has started since.
161    pub async fn call(&self, req: Req) -> Result<Res, Err> {
162        let pl = self.clone();
163        pl.api.call(pl.idx, req).await
164    }
165
166    #[inline]
167    /// Returns an owned future that waits for readiness and calls the service.
168    ///
169    /// The returned future does not borrow this binding and can be moved between
170    /// local tasks. The readiness check is skipped if the last pipeline readiness
171    /// check succeeded and no call has started since; this is decided when the
172    /// future is first polled.
173    pub fn call_static(&self, req: Req) -> PipelineCall<Req, Res, Err> {
174        PipelineCall::new(self.clone(), req)
175    }
176}
177
178impl<Req, Res, Err> Drop for PipelineBinding<Req, Res, Err> {
179    #[inline]
180    fn drop(&mut self) {
181        self.api.unreg(self.idx);
182    }
183}
184
185impl<Req, Res, Err> Clone for PipelineBinding<Req, Res, Err> {
186    fn clone(&self) -> Self {
187        Self {
188            idx: self.api.reg(),
189            api: self.api.clone(),
190        }
191    }
192}
193
194#[must_use = "futures do nothing unless polled"]
195/// An owned future for a pipeline service call.
196///
197/// Created by [`Pipeline::call_static`] and [`PipelineBinding::call_static`]. The
198/// future keeps the pipeline alive. It does not check whether the pipeline is shut
199/// down; the request is passed to the service as usual.
200pub struct PipelineCall<Req, Res, Err> {
201    // `fut` borrows from `pl`, so it must be declared (and dropped) first
202    fut: BoxFuture<'static, Result<Res, Err>>,
203    #[allow(dead_code)]
204    pl: PipelineBinding<Req, Res, Err>,
205}
206
207impl<Req, Res, Err> PipelineCall<Req, Res, Err> {
208    #[allow(clippy::missing_transmute_annotations)]
209    fn new(pl: PipelineBinding<Req, Res, Err>, req: Req) -> Self {
210        // SAFETY: `fut` borrows from `pl.api` (`Rc`-allocated, never moves).
211        // `fut` is declared before `pl` in `PipelineCall`, so it is dropped first.
212        PipelineCall {
213            fut: unsafe { std::mem::transmute(pl.api.call(pl.idx, req)) },
214            pl,
215        }
216    }
217}
218
219impl<Req, Res, Err> future::Future for PipelineCall<Req, Res, Err> {
220    type Output = Result<Res, Err>;
221
222    #[inline]
223    fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
224        Pin::new(&mut self.as_mut().fut).poll(cx)
225    }
226}
227
228impl<Req, Res, Err> fmt::Debug for Pipeline<Req, Res, Err> {
229    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
230        f.debug_struct("Pipeline").finish()
231    }
232}
233
234impl<Req, Res, Err> fmt::Debug for PipelineBinding<Req, Res, Err> {
235    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
236        f.debug_struct("PipelineBinding")
237            .field("idx", &self.idx)
238            .finish()
239    }
240}
241
242impl<Req, Res, Err> fmt::Debug for PipelineCall<Req, Res, Err> {
243    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
244        f.debug_struct("PipelineCall").finish()
245    }
246}