1#![deny(clippy::pedantic)]
8#![allow(
9 clippy::cast_possible_truncation,
10 clippy::missing_fields_in_debug,
11 clippy::missing_errors_doc,
12 clippy::missing_panics_doc,
13 clippy::must_use_candidate,
14 clippy::type_complexity,
15 clippy::unused_async,
16 clippy::unused_async_trait_impl
17)]
18use std::rc::Rc;
19
20mod and_then;
21mod apply;
22pub mod boxed;
23pub mod cfg;
24mod chain;
25mod ctx;
26mod fn_ready;
27mod fn_service;
28mod fn_shutdown;
29mod macros;
30mod map;
31mod map_err;
32mod map_init_err;
33mod map_state;
34mod middleware;
35pub mod state;
36mod then;
37mod util;
38
39pub mod pipeline;
40mod pl_factory;
41mod pl_inner;
42mod pl_state;
43
44pub use crate::apply::{apply_fn, apply_fn_factory};
45pub use crate::chain::{ServiceChain, ServiceChainFactory, factory, service};
46pub use crate::ctx::Ctx;
47pub use crate::fn_service::{fn_factory, fn_service, fn_service_st};
48pub use crate::map_state::{map_state, map_state_factory};
49pub use crate::middleware::{Identity, Middleware, Stack, apply, fn_layer};
50pub use crate::pipeline::Pipeline;
51pub use crate::state::{RequestState, State};
52
53#[allow(unused_variables)]
54pub trait Service<St, Req> {
114 type Res;
116
117 type Error;
119
120 async fn call(&self, req: Req, ctx: Ctx<'_, Self, St>) -> Result<Self::Res, Self::Error>;
126
127 #[inline]
128 async fn ready(&self, ctx: Ctx<'_, Self, St>) -> Result<(), Self::Error> {
136 Ok(())
137 }
138
139 #[inline]
140 async fn shutdown(&self, ctx: Ctx<'_, Self, St>) {}
144
145 #[inline]
146 fn map<F, Res>(self, f: F) -> ServiceChain<dev::Map<F, Self, Res>, St, Req>
154 where
155 Self: Sized,
156 F: Fn(Self::Res) -> Res,
157 {
158 service(dev::Map::new(f, self))
159 }
160
161 #[inline]
162 fn map_err<F, E>(self, f: F) -> ServiceChain<dev::MapErr<F, Self, E>, St, Req>
170 where
171 Self: Sized,
172 F: Fn(Self::Error) -> E,
173 {
174 service(dev::MapErr::new(f, self))
175 }
176
177 #[inline]
178 fn and_then<Next, F>(self, f: F) -> ServiceChain<dev::AndThen<Self, Next>, St, Req>
183 where
184 Self: Sized,
185 Next: Service<St, Self::Res, Error = Self::Error>,
186 F: IntoService<Next, St, Self::Res>,
187 {
188 service(dev::AndThen::new(self, f.into_service()))
189 }
190
191 #[inline]
192 fn pipeline(self, st: St) -> Pipeline<Req, Self::Res, Self::Error>
194 where
195 Self: Sized + 'static,
196 St: 'static,
197 Req: 'static,
198 {
199 Pipeline::new(st, self)
200 }
201}
202
203pub trait ServiceFactory<St, Req> {
214 type Res;
216
217 type Error;
219
220 type Service: Service<St, Req, Res = Self::Res, Error = Self::Error>;
222
223 type InitError;
225
226 async fn create(&self, st: &St) -> Result<Self::Service, Self::InitError>;
228
229 #[inline]
230 async fn pipeline(
232 &self,
233 st: St,
234 ) -> Result<Pipeline<Req, Self::Res, Self::Error>, Self::InitError>
235 where
236 Self: 'static,
237 St: 'static,
238 Req: 'static,
239 {
240 let svc = self.create(&st).await?;
241 Ok(Pipeline::new(st, svc))
242 }
243
244 #[inline]
245 fn map<F, Res>(self, f: F) -> ServiceChainFactory<dev::MapFactory<F, Self, Res>, St, Req>
247 where
248 Self: Sized,
249 F: Fn(Self::Res) -> Res + Clone,
250 {
251 factory(dev::MapFactory::new(f, self))
252 }
253
254 #[inline]
255 fn map_err<F, E>(self, f: F) -> ServiceChainFactory<dev::MapErrFactory<F, Self, E>, St, Req>
257 where
258 Self: Sized,
259 F: Fn(Self::Error) -> E + Clone,
260 {
261 factory(dev::MapErrFactory::new(f, self))
262 }
263
264 #[inline]
265 fn map_init_err<F, E>(self, f: F) -> ServiceChainFactory<dev::MapInitErr<F, Self, E>, St, Req>
268 where
269 Self: Sized,
270 F: Fn(Self::InitError) -> E + Clone,
271 {
272 factory(dev::MapInitErr::new(f, self))
273 }
274
275 fn and_then<U, F>(self, f: F) -> ServiceChainFactory<dev::AndThenFactory<Self, U>, St, Req>
281 where
282 Self: Sized,
283 U: ServiceFactory<St, Self::Res, Error = Self::Error, InitError = Self::InitError>,
284 F: IntoServiceFactory<U, St, Self::Res>,
285 {
286 factory(dev::AndThenFactory::new(self, f.into_factory()))
287 }
288
289 fn boxed(self) -> boxed::BoxServiceFactory<St, Req, Self::Res, Self::Error, Self::InitError>
291 where
292 St: 'static,
293 Req: 'static,
294 Self: Sized + 'static,
295 {
296 boxed::factory(self)
297 }
298}
299
300impl<S, St, Req> Service<St, Req> for &S
301where
302 S: Service<St, Req>,
303{
304 type Res = S::Res;
305 type Error = S::Error;
306
307 #[inline]
308 async fn ready(&self, ctx: Ctx<'_, Self, St>) -> Result<(), S::Error> {
309 ctx.ready(&**self).await
310 }
311
312 #[inline]
313 async fn call(&self, req: Req, ctx: Ctx<'_, Self, St>) -> Result<S::Res, S::Error> {
314 ctx.call_nowait(&**self, req).await
315 }
316
317 #[inline]
318 async fn shutdown(&self, ctx: Ctx<'_, Self, St>) {
319 ctx.shutdown(&**self).await;
320 }
321}
322
323impl<S, St, Req> Service<St, Req> for Box<S>
324where
325 S: Service<St, Req>,
326{
327 type Res = S::Res;
328 type Error = S::Error;
329
330 #[inline]
331 async fn ready(&self, ctx: Ctx<'_, Self, St>) -> Result<(), S::Error> {
332 ctx.ready(&**self).await
333 }
334
335 #[inline]
336 async fn call(&self, req: Req, ctx: Ctx<'_, Self, St>) -> Result<S::Res, S::Error> {
337 ctx.call_nowait(&**self, req).await
338 }
339
340 #[inline]
341 async fn shutdown(&self, ctx: Ctx<'_, Self, St>) {
342 ctx.shutdown(&**self).await;
343 }
344}
345
346impl<S, St, Req> Service<St, Req> for Rc<S>
347where
348 S: Service<St, Req>,
349{
350 type Res = S::Res;
351 type Error = S::Error;
352
353 #[inline]
354 async fn ready(&self, ctx: Ctx<'_, Self, St>) -> Result<(), S::Error> {
355 ctx.ready(&**self).await
356 }
357
358 #[inline]
359 async fn call(&self, req: Req, ctx: Ctx<'_, Self, St>) -> Result<S::Res, S::Error> {
360 ctx.call_nowait(&**self, req).await
361 }
362
363 #[inline]
364 async fn shutdown(&self, ctx: Ctx<'_, Self, St>) {
365 ctx.shutdown(&**self).await;
366 }
367}
368
369impl<Sf, St, Req> ServiceFactory<St, Req> for Rc<Sf>
370where
371 Sf: ServiceFactory<St, Req>,
372{
373 type Res = Sf::Res;
374 type Error = Sf::Error;
375 type Service = Sf::Service;
376 type InitError = Sf::InitError;
377
378 async fn create(&self, st: &St) -> Result<Self::Service, Self::InitError> {
379 self.as_ref().create(st).await
380 }
381}
382
383pub trait ServiceCaller<Req, Res, Err> {
385 async fn call_service(&self, req: Req) -> Result<Res, Err>;
387}
388
389pub trait IntoService<S, St, Req>
391where
392 S: Service<St, Req>,
393{
394 fn into_service(self) -> S;
396}
397
398pub trait IntoServiceFactory<Sf, St, Req>
400where
401 Sf: ServiceFactory<St, Req>,
402{
403 fn into_factory(self) -> Sf;
405}
406
407impl<S, St, Req> IntoService<S, St, Req> for S
408where
409 S: Service<St, Req>,
410{
411 #[inline]
412 fn into_service(self) -> S {
413 self
414 }
415}
416
417impl<Sf, St, Req> IntoServiceFactory<Sf, St, Req> for Sf
418where
419 Sf: ServiceFactory<St, Req>,
420{
421 #[inline]
422 fn into_factory(self) -> Sf {
423 self
424 }
425}
426
427pub mod dev {
429 pub use crate::and_then::{AndThen, AndThenFactory};
430 pub use crate::apply::{Apply, ApplyCtx, ApplyFactory};
431 pub use crate::chain::{ServiceChain, ServiceChainFactory};
432 pub use crate::fn_ready::FnReadiness;
433 pub use crate::fn_service::{FnFactory, FnService, FnServiceSt, FnServiceStFactory};
434 pub use crate::fn_shutdown::FnShutdown;
435 pub use crate::map::{Map, MapFactory};
436 pub use crate::map_err::{MapErr, MapErrFactory};
437 pub use crate::map_init_err::MapInitErr;
438 pub use crate::map_state::{MapState, MapStateFactory};
439 pub use crate::middleware::{ApplyMiddleware, FnMiddleware};
440 pub use crate::then::{Then, ThenFactory};
441}
442
443#[cfg(test)]
444mod tests {
445 use std::{cell::Cell, task::Poll};
446
447 use ntex::util::lazy;
448
449 use super::*;
450 use crate::dev::FnServiceSt;
451 use crate::pipeline::PipelineFactory;
452
453 #[derive(Clone, Debug, Default)]
454 struct Srv(Rc<Cell<usize>>);
455
456 impl Service<usize, usize> for Srv {
457 type Res = usize;
458 type Error = &'static str;
459
460 async fn ready(&self, ctx: Ctx<'_, Self, usize>) -> Result<(), Self::Error> {
461 assert!(ctx.poll_once(|_| true));
463 self.0.set(self.0.get() + 1);
464 Ok(())
465 }
466
467 async fn call(&self, req: usize, ctx: Ctx<'_, Self, usize>) -> Result<usize, Self::Error> {
468 let one = ctx.poll_once(|_| 1);
469 let two = ctx.poll_fn(|_| Poll::Ready(2)).await;
470 assert_eq!(one + two, 3);
471
472 if req == 0 { Err("zero") } else { Ok(req + *ctx) }
473 }
474
475 async fn shutdown(&self, _: Ctx<'_, Self, usize>) {
476 self.0.set(self.0.get() + 100);
477 }
478 }
479
480 struct Fwd {
481 inner: Srv,
482 pl: Pipeline<usize, usize, &'static str>,
483 }
484
485 impl Service<usize, usize> for Fwd {
486 type Res = usize;
487 type Error = usize;
488
489 crate::forward_ready!(usize, inner, str::len);
490
491 async fn call(&self, req: usize, ctx: Ctx<'_, Self, usize>) -> Result<usize, usize> {
492 let res = ctx.call(&self.inner, req).await.map_err(str::len)?;
493 self.pl.call(res).await.map_err(str::len)
494 }
495
496 crate::forward_shutdown!(usize, inner);
497 }
498
499 struct FwdPl {
500 pl: Pipeline<usize, usize, &'static str>,
501 }
502
503 impl Service<usize, usize> for FwdPl {
504 type Res = usize;
505 type Error = &'static str;
506
507 crate::forward_pl_ready!(usize, pl);
508 crate::forward_pl_shutdown!(usize, pl);
509
510 async fn call(&self, req: usize, _: Ctx<'_, Self, usize>) -> Result<usize, &'static str> {
511 self.pl.call(req).await
512 }
513 }
514
515 struct FwdPlErr {
516 pl: Pipeline<usize, usize, &'static str>,
517 }
518
519 impl Service<usize, usize> for FwdPlErr {
520 type Res = usize;
521 type Error = usize;
522
523 crate::forward_pl_ready!(usize, pl, str::len);
524
525 async fn call(&self, req: usize, _: Ctx<'_, Self, usize>) -> Result<usize, usize> {
526 self.pl.call(req).await.map_err(str::len)
527 }
528 }
529
530 #[ntex::test]
531 async fn service_wrappers() {
532 let cnt = Rc::new(Cell::new(0));
533
534 let srv: &'static Srv = Box::leak(Box::new(Srv(cnt.clone())));
535 let pl = Pipeline::new(1, srv);
536 assert_eq!(pl.call(1).await, Ok(2));
537 assert_eq!(pl.call(0).await, Err("zero"));
538 pl.shutdown().await;
539 assert_eq!(cnt.get(), 102);
540
541 let pl = Pipeline::new(2, Box::new(Srv(cnt.clone())));
542 assert_eq!(pl.call(1).await, Ok(3));
543 pl.shutdown().await;
544 assert_eq!(cnt.get(), 203);
545
546 let pl = Pipeline::new(3, Rc::new(Srv(cnt.clone())));
547 assert_eq!(pl.call(1).await, Ok(4));
548 pl.shutdown().await;
549 assert_eq!(cnt.get(), 304);
550
551 let pl = Pipeline::new(
552 1,
553 Fwd {
554 inner: Srv(cnt.clone()),
555 pl: Pipeline::new(10, Srv(cnt.clone())),
556 },
557 );
558 assert_eq!(pl.call(1).await, Ok(12));
559 assert_eq!(pl.call(0).await, Err(4));
560 pl.shutdown().await;
561 assert_eq!(cnt.get(), 409);
562
563 let pl = Pipeline::new(
564 1,
565 FwdPl {
566 pl: Pipeline::new(10, Srv(cnt.clone())),
567 },
568 );
569 assert_eq!(pl.call(1).await, Ok(11));
570 pl.shutdown().await;
571 assert_eq!(cnt.get(), 510);
573
574 let pl = Pipeline::new(
575 1,
576 FwdPlErr {
577 pl: Pipeline::new(10, Srv(cnt.clone())),
578 },
579 );
580 assert_eq!(pl.call(0).await, Err(4));
581 pl.shutdown().await;
582 assert_eq!(cnt.get(), 511);
583 }
584
585 #[ntex::test]
586 async fn pipeline_calls() {
587 let cnt = Rc::new(Cell::new(0));
588 let pl = Srv(cnt.clone()).pipeline(1);
589
590 assert_eq!(pl.call_static(1).await, Ok(2));
591 assert_eq!(cnt.get(), 1);
592
593 assert_eq!(pl.ready().await, Ok(()));
595 assert_eq!(cnt.get(), 2);
596 assert_eq!(pl.call(2).await, Ok(3));
597 assert_eq!(cnt.get(), 2);
598 assert_eq!(ServiceCaller::call_service(&pl, 3).await, Ok(4));
599 assert_eq!(cnt.get(), 3);
600
601 assert_eq!(lazy(|cx| pl.poll_ready(cx)).await, Poll::Ready(Ok(())));
602 assert_eq!(cnt.get(), 4);
603 assert_eq!(pl.call_static(3).await, Ok(4));
604 assert_eq!(pl.call_static(3).await, Ok(4));
605 assert_eq!(cnt.get(), 5);
606
607 let b = pl.bind();
608 assert!(format!("{b:?}").contains("PipelineBinding"));
609 assert_eq!(b.call_static(4).await, Ok(5));
610 assert_eq!(cnt.get(), 6);
611 assert_eq!(b.ready().await, Ok(()));
612 assert_eq!(cnt.get(), 7);
613 assert_eq!(pl.call(4).await, Ok(5));
614 assert_eq!(b.call(4).await, Ok(5));
615 assert_eq!(cnt.get(), 8);
616
617 assert_eq!(pl.ready().await, Ok(()));
619 let fut1 = pl.call_static(1);
620 let fut2 = pl.call_static(2);
621 assert_eq!(fut2.await, Ok(3));
622 assert_eq!(cnt.get(), 9);
623 assert_eq!(fut1.await, Ok(2));
624 assert_eq!(cnt.get(), 10);
625
626 assert_eq!(pl.ready().await, Ok(()));
628 assert_eq!(cnt.get(), 11);
629 pl.shutdown().await;
630 assert_eq!(pl.call(1).await, Ok(2));
631 assert_eq!(cnt.get(), 112);
632
633 let svc = apply_fn(
634 Srv(cnt.clone()),
635 async |req: usize, svc: &dev::ApplyCtx<'_, Srv, usize, usize>| {
636 svc.call_service(req * 2).await
637 },
638 );
639 let pl = Pipeline::new(1, svc);
640 assert_eq!(pl.call(2).await, Ok(5));
641 assert_eq!(cnt.get(), 115);
642 }
643
644 #[ntex::test]
645 async fn fn_conversions() {
646 let pl = Pipeline::new(5, async |st: &usize, req: usize| Ok::<_, ()>(st + req));
647 assert_eq!(pl.call(1).await, Ok(6));
648
649 let _: FnServiceSt<_, usize, usize, usize, ()> =
650 IntoService::into_service(async |st: &usize, req: usize| Ok::<_, ()>(st + req));
651
652 let f = factory(async |st: &usize| Ok::<_, ()>(Srv(Rc::new(Cell::new(*st)))));
653 let pl = f.pipeline(1).await.unwrap();
654 assert_eq!(pl.call(1).await, Ok(2));
655 }
656
657 #[ntex::test]
658 async fn factory_combinators() {
659 let cnt = Rc::new(Cell::new(0));
660 let c = cnt.clone();
661 let f = fn_factory(async move |st: &usize| {
662 if *st == 0 {
663 Err(())
664 } else {
665 Ok::<_, ()>(Srv(c.clone()))
666 }
667 });
668
669 let rc = Rc::new(f.clone());
670 let pl = rc.pipeline(1).await.unwrap();
671 assert_eq!(pl.call(1).await, Ok(2));
672 assert!(rc.pipeline(0).await.is_err());
673
674 let f2 = f.clone().map_init_err(|()| "init");
675 assert_eq!(f2.create(&0).await.err(), Some("init"));
676
677 let f2 = f.clone().and_then(f.clone());
678 let pl = f2.pipeline(1).await.unwrap();
679 assert_eq!(pl.call(1).await, Ok(3));
680
681 let f2 = factory(f.clone()).map(|r| r * 10).map_err(|_| ());
682 let pl = f2.pipeline(1).await.unwrap();
683 assert_eq!(pl.call(1).await, Ok(20));
684 assert_eq!(pl.call(0).await, Err(()));
685
686 let pf = PipelineFactory::new(f.clone());
687 let pf2 = pf.clone();
688 assert!(format!("{pf2:?}").contains("PipelineFactory"));
689 let pl = pf2.create(2).await.unwrap();
690 assert_eq!(pl.call(1).await, Ok(3));
691 assert!(pf.create(0).await.is_err());
692 }
693
694 #[ntex::test]
695 async fn boxed_and_debug() {
696 let cnt = Rc::new(Cell::new(0));
697
698 let svc = boxed::service(Srv(cnt.clone()));
699 let pl = Pipeline::new(1, svc.clone());
700 assert_eq!(pl.call(1).await, Ok(2));
701
702 let f = boxed::factory(fn_factory(async |_: &usize| Ok::<_, ()>(Srv::default())));
703 let pl = f.clone().pipeline(1).await.unwrap();
704 assert_eq!(pl.call(1).await, Ok(2));
705
706 let s = format!("{:?}", service(Srv::default()).map(|r| r).map_err(|e| e));
707 assert!(s.contains("Map") && s.contains("MapErr"));
708 }
709
710 #[test]
711 fn request_state() {
712 let st = State {
713 req: 1,
714 state: "st",
715 };
716 assert_eq!(st.unpack(), ("st", 1));
717 assert_eq!(("st", 2).unpack(), ("st", 2));
718 }
719}