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
9pub struct Pipeline<Req, Res, Err> {
14 api: PipelineApi<Req, Res, Err>,
15}
16
17pub 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 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 pub async fn ready(&self) -> Result<(), Err> {
50 future::poll_fn(|cx| self.api.poll_ready(cx)).await
51 }
52
53 #[inline]
54 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 pub fn call_static(&self, req: Req) -> PipelineCall<Req, Res, Err> {
71 PipelineCall::new(self.bind(), req)
72 }
73
74 #[inline]
75 pub fn poll_ready(&self, cx: &mut Context<'_>) -> Poll<Result<(), Err>> {
80 self.api.poll_ready(cx)
81 }
82
83 #[inline]
84 pub fn poll_shutdown(&self, cx: &mut Context<'_>) -> Poll<()> {
86 self.api.poll_shutdown(cx)
87 }
88
89 #[inline]
90 pub fn is_shutdown(&self) -> bool {
92 self.api.is_shutdown()
93 }
94
95 #[inline]
96 pub async fn shutdown(&self) {
98 future::poll_fn(|cx| self.api.poll_shutdown(cx)).await;
99 }
100
101 #[inline]
102 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 pub async fn ready(&self) -> Result<(), Err> {
153 self.api.ready(self.idx).await
154 }
155
156 #[inline]
157 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 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"]
195pub struct PipelineCall<Req, Res, Err> {
201 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 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}