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#[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 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 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 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 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 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}