Skip to main content

ntex/http/h1/
service.rs

1use crate::error::IntoFailure;
2use crate::http::error::{DispatchError, ResponseError};
3use crate::http::{Request, Response, config::DispatcherConfig};
4use crate::io::{Filter, Io, types};
5use crate::service::pipeline::{Pipeline, PipelineFactory};
6use crate::service::{Ctx, IntoServiceFactory, RequestState, Service, ServiceFactory};
7
8use super::control::{Control, ControlAck, ControlResult};
9use super::dispatcher::Dispatcher;
10
11/// Server transport service that dispatches HTTP/1 requests.
12///
13/// Construct this service with [`HttpService::h1`](crate::http::HttpService::h1).
14/// Request handling is delegated to the configured application service, while
15/// connection lifecycle events get their default action unless a control
16/// service is provided with [`control`](Self::control).
17#[derive(derive_more::Debug)]
18#[debug("H1Service")]
19pub struct H1Service<F, Req: RequestState<Io<F>>, Err> {
20    sf: crate::http::HttpPipeline<Req::State, Err>,
21    /// Without a control service the default action is applied to every event.
22    ctl: Option<crate::http::Ctl1Pipeline<Req::State, F, Err>>,
23    config: DispatcherConfig,
24}
25
26impl<F, Req, Err> H1Service<F, Req, Err>
27where
28    F: Filter,
29    Req: RequestState<Io<F>>,
30    Req::State: Clone,
31    Err: ResponseError + 'static,
32{
33    /// Create new `H1Service` instance.
34    pub(crate) fn new<Sf>(sf: impl IntoServiceFactory<Sf, Req::State, Request>) -> Self
35    where
36        Sf: ServiceFactory<Req::State, Request, Error = Err> + 'static,
37        Sf::Res: Into<Response>,
38        Sf::InitError: IntoFailure,
39    {
40        H1Service {
41            sf: PipelineFactory::new(
42                sf.into_factory()
43                    .map(Into::into)
44                    .map_init_err(|e| DispatchError::Control(e.fail())),
45            ),
46            ctl: None,
47            config: DispatcherConfig::default(),
48        }
49    }
50}
51
52impl<F, Req, Err> H1Service<F, Req, Err>
53where
54    F: Filter,
55    Req: RequestState<Io<F>>,
56    Req::State: Clone,
57    Err: ResponseError + 'static,
58{
59    #[must_use]
60    /// Provides the HTTP/1 control service.
61    ///
62    /// The control service receives connection, request, expectation, upgrade,
63    /// and disconnect events. Service and protocol errors are reported by the
64    /// disconnect event. Returning [`Control::ack`] applies the
65    /// default action for each event.
66    pub fn control<I, Sf>(self, ctl: I) -> Self
67    where
68        I: IntoServiceFactory<Sf, Req::State, Control<F, Err>>,
69        Sf: ServiceFactory<Req::State, Control<F, Err>, Res = ControlAck<F>> + 'static,
70        Sf::Error: IntoFailure,
71        Sf::InitError: IntoFailure,
72    {
73        H1Service {
74            sf: self.sf,
75            ctl: Some(PipelineFactory::new(
76                ctl.into_factory()
77                    .map_err(|e| DispatchError::Service(e.fail()))
78                    .map_init_err(|e| DispatchError::Control(e.fail())),
79            )),
80            config: self.config,
81        }
82    }
83}
84
85impl<St, F, Req, Err> Service<St, Req> for H1Service<F, Req, Err>
86where
87    F: Filter,
88    Req: RequestState<Io<F>>,
89    Req::State: Clone,
90    Err: ResponseError + 'static,
91{
92    type Res = ();
93    type Error = DispatchError;
94
95    async fn call(&self, req: Req, _: Ctx<'_, Self, St>) -> Result<(), Self::Error> {
96        let (st, io) = req.unpack();
97
98        let svc = self.sf.create(st.clone()).await?;
99        let ctl = match &self.ctl {
100            Some(ctl) => Some(ctl.create(st).await?),
101            None => None,
102        };
103
104        let id = self.config.next_id();
105        let ioref = io.get_ref();
106        let (_guard, inflight) = self.config.insert_io(&ioref);
107
108        log::trace!(
109            "{}: New http1 connection {id}, peer address {:?}, inflight: {}",
110            io.tag(),
111            io.query::<types::PeerAddr>().get(),
112            inflight
113        );
114
115        handle_io(id, io, svc, ctl, self.config.clone()).await
116    }
117
118    async fn ready(&self, _: Ctx<'_, Self, St>) -> Result<(), Self::Error> {
119        Ok(())
120    }
121
122    async fn shutdown(&self, _: crate::Ctx<'_, Self, St>) {
123        // check inflight connections
124        let inflight = self.config.shutdown();
125        if inflight != 0 {
126            log::trace!("Shutting down service, in-flight connections: {inflight}");
127
128            self.config.wait_shutdown().await;
129            log::trace!("Shutting down is complected");
130        }
131    }
132}
133
134pub(crate) async fn handle_io<F, Err>(
135    id: usize,
136    io: Io<F>,
137    svc: Pipeline<Request, Response, Err>,
138    ctl: Option<Pipeline<Control<F, Err>, ControlAck<F>, DispatchError>>,
139    config: DispatcherConfig,
140) -> Result<(), DispatchError>
141where
142    F: Filter,
143    Err: ResponseError + 'static,
144{
145    // Notify control service
146    let io = if let Some(ctl) = &ctl {
147        let ack = ctl.call(Control::connect(id, io)).await?;
148        let ControlResult::Connect(io) = ack.result else {
149            unreachable!();
150        };
151        io
152    } else {
153        io
154    };
155
156    Dispatcher::new(id, io, svc, ctl, config).await
157}