1use super::Upgraded;
6use crate::opentelemetry::version_as_protocol_version;
7use rama_core::error::ErrorExt as _;
8use rama_core::error_sink::ErrorSink;
9use rama_core::extensions::ExtensionsRef;
10use rama_core::layer::ConsumeErr;
11use rama_core::rt::Executor;
12use rama_core::telemetry::tracing::{self, Instrument};
13use rama_core::{Service, extensions::Extensions, matcher::Matcher, service::BoxService};
14use rama_http_types::Request;
15use rama_utils::macros::define_inner_service_accessors;
16use std::{convert::Infallible, fmt, sync::Arc};
17
18pub struct UpgradeService<S, O> {
21 handlers: Vec<Arc<UpgradeHandler<O>>>,
22 inner: S,
23 exec: Executor,
24 error_sink: Arc<dyn ErrorSink>,
25}
26
27#[derive(Clone, Debug)]
28pub struct UpgradeResponse<I, O> {
29 pub response: O,
31 pub request: I,
33 pub extensions: Extensions,
36}
37
38pub struct UpgradeHandler<O> {
40 matcher: Box<dyn Matcher<Request>>,
41 responder: BoxService<Request, UpgradeResponse<Request, O>, O>,
42 handler: BoxService<Upgraded, (), Infallible>,
45 _phantom: std::marker::PhantomData<fn(O) -> ()>,
46}
47
48impl<O> UpgradeHandler<O> {
49 pub(crate) fn new<M, R, H, Sink>(matcher: M, responder: R, handler: H, sink: Sink) -> Self
52 where
53 M: Matcher<Request>,
54 R: Service<Request, Output = UpgradeResponse<Request, O>, Error = O> + Clone,
55 H: Service<Upgraded, Output = ()> + Clone,
56 Sink: ErrorSink<H::Error>,
57 {
58 let sink = Arc::new(sink);
59 let handler = ConsumeErr::new(handler, move |err| sink.sink_error(err)).boxed();
62 Self {
63 matcher: Box::new(matcher),
64 responder: responder.boxed(),
65 handler,
66 _phantom: std::marker::PhantomData,
67 }
68 }
69}
70
71impl<S, O> UpgradeService<S, O> {
72 pub fn new(
74 handlers: Vec<Arc<UpgradeHandler<O>>>,
75 inner: S,
76 exec: Executor,
77 error_sink: Arc<dyn ErrorSink>,
78 ) -> Self {
79 Self {
80 handlers,
81 inner,
82 exec,
83 error_sink,
84 }
85 }
86
87 define_inner_service_accessors!();
88}
89
90impl<S, O> fmt::Debug for UpgradeService<S, O>
91where
92 S: fmt::Debug,
93{
94 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
95 f.debug_struct("UpgradeService")
96 .field("handlers", &self.handlers)
97 .field("inner", &self.inner)
98 .field("exec", &self.exec)
99 .finish()
100 }
101}
102
103impl<S, O> Clone for UpgradeService<S, O>
104where
105 S: Clone,
106{
107 fn clone(&self) -> Self {
108 Self {
109 handlers: self.handlers.clone(),
110 inner: self.inner.clone(),
111 exec: self.exec.clone(),
112 error_sink: self.error_sink.clone(),
113 }
114 }
115}
116
117impl<S, O, E> Service<Request> for UpgradeService<S, O>
118where
119 S: Service<Request, Output = O, Error = E>,
120 O: Send + Sync + 'static,
121 E: Send + Sync + 'static,
122{
123 type Output = O;
124 type Error = E;
125
126 async fn serve(&self, req: Request) -> Result<Self::Output, Self::Error> {
127 for handler in &self.handlers {
128 let ext = Extensions::new();
129 if !handler.matcher.matches(Some(&ext), &req) {
130 continue;
131 }
132 req.extensions().extend(&ext);
133
134 return match handler.responder.serve(req).await {
135 Ok(UpgradeResponse {
136 response,
137 request,
138 extensions,
139 }) => {
140 let handler = handler.handler.clone();
141 let error_sink = self.error_sink.clone();
142
143 let span = tracing::trace_root_span!(
144 "upgrade::serve",
145 otel.kind = "server",
146 http.request.method = %request.method().as_str(),
147 url.full = %request.request_uri(),
148 url.path = %request.uri().path_or_root().as_ref(),
149 url.query = %request.uri().query_or_empty().as_ref(),
150 url.scheme = %request.uri().scheme_str().unwrap_or_default(),
151 network.protocol.name = "http",
152 network.protocol.version = version_as_protocol_version(request.version()),
153 );
154
155 self.exec.spawn_task(
156 async move {
157 match crate::io::upgrade::handle_upgrade(request).await {
158 Ok(upgraded) => {
159 upgraded.extensions().extend(&extensions);
160 _ = handler.serve(upgraded).await;
164 }
165 Err(err) => {
166 error_sink.sink_error(
169 err.context("http upgrade failed before handler"),
170 );
171 }
172 }
173 }
174 .instrument(span),
175 );
176 Ok(response)
177 }
178 Err(e) => Ok(e),
179 };
180 }
181
182 self.inner.serve(req).await
183 }
184}
185
186impl<O> fmt::Debug for UpgradeHandler<O> {
187 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
188 f.debug_struct("UpgradeHandler").finish()
189 }
190}
191
192#[cfg(test)]
193mod tests {
194 use super::*;
195 use crate::io::upgrade::{Upgraded, pending};
196 use crate::layer::upgrade::UpgradeLayer;
197 use rama_core::Layer;
198 use rama_core::ServiceInput;
199 use rama_core::bytes::Bytes;
200 use rama_core::error::{BoxError, BoxErrorExt as _};
201 use rama_core::service::service_fn;
202 use rama_http_types::{Body, Response};
203 use std::convert::Infallible;
204 use std::time::Duration;
205 use tokio::sync::mpsc;
206 use tokio_test::io::Builder;
207
208 #[tokio::test]
211 async fn upgrade_handler_error_is_routed_to_sink() {
212 let (tx, mut rx) = mpsc::unbounded_channel::<String>();
214
215 let (pending_upgrade, on_upgrade) = pending();
217 let req = Request::new(Body::empty());
218 req.extensions().insert(on_upgrade);
219
220 let responder = service_fn(|req: Request| async move {
223 Ok::<_, Response>(UpgradeResponse {
224 response: Response::new(Body::empty()),
225 request: req,
226 extensions: Extensions::new(),
227 })
228 });
229
230 let handler = service_fn(|_upgraded: Upgraded| async move {
232 Err::<(), BoxError>(BoxError::from_static_str("handler boom"))
233 });
234
235 let inner =
237 service_fn(
238 |_req: Request| async move { Ok::<_, Infallible>(Response::new(Body::empty())) },
239 );
240
241 let svc = UpgradeLayer::new_with_error_sink(
244 Executor::default(),
245 true,
246 responder,
247 handler,
248 move |err: BoxError| {
249 _ = tx.send(format!("{err:?}"));
250 },
251 )
252 .into_layer(inner);
253
254 let _resp = svc.serve(req).await.expect("upgrade match -> Ok(response)");
256
257 let upgraded = Upgraded::new(ServiceInput::new(Builder::default().build()), Bytes::new());
259 pending_upgrade.fulfill(upgraded);
260
261 let reported = tokio::time::timeout(Duration::from_secs(5), rx.recv())
263 .await
264 .expect("sink should be called within timeout")
265 .expect("sink channel should yield the error");
266 assert!(
267 reported.contains("handler boom"),
268 "unexpected sink message: {reported}"
269 );
270 }
271}