Skip to main content

rama_http/layer/upgrade/
service.rs

1//! upgrade service to handle branching into http upgrade services
2//!
3//! See [`UpgradeService`] for more details.
4
5use 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
18/// Upgrade service can be used to handle the possibility of upgrading a request,
19/// after which it will pass down the transport RW to the attached upgrade service.
20pub 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    /// Response that should be returned
30    pub response: O,
31    /// Request that caused this upgrade
32    pub request: I,
33    /// Extensions which will be applied to the [`Upgraded`] io
34    /// if the upgrade was successful
35    pub extensions: Extensions,
36}
37
38/// UpgradeHandler is a helper struct used internally to create an upgrade service.
39pub struct UpgradeHandler<O> {
40    matcher: Box<dyn Matcher<Request>>,
41    responder: BoxService<Request, UpgradeResponse<Request, O>, O>,
42    // The handler's own error (any type `E`) is consumed by its [`ErrorSink`]
43    // inside this boxed unit, so nothing remains to propagate (`Infallible`).
44    handler: BoxService<Upgraded, (), Infallible>,
45    _phantom: std::marker::PhantomData<fn(O) -> ()>,
46}
47
48impl<O> UpgradeHandler<O> {
49    /// Create a new upgrade handler whose own errors (of any type `E`) are
50    /// routed to the given [`ErrorSink`].
51    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        // Consume the handler's error in place via its sink, so the boxed
60        // handler is `Infallible` regardless of the handler's error type.
61        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    /// Create a new [`UpgradeService`].
73    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                                    // The handler's own error (if any) was already
161                                    // consumed by its per-handler [`ErrorSink`]; the
162                                    // boxed handler is `Infallible` here.
163                                    _ = handler.serve(upgraded).await;
164                                }
165                                Err(err) => {
166                                    // The HTTP upgrade itself failed (before the handler
167                                    // ran): route it to the layer's upgrade error sink.
168                                    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    // Regression for #1014: a failing upgrade handler must hand its error to its
209    // per-handler [`ErrorSink`] instead of being silently swallowed.
210    #[tokio::test]
211    async fn upgrade_handler_error_is_routed_to_sink() {
212        // mpsc so the (sync) sink can report out of the detached upgrade task.
213        let (tx, mut rx) = mpsc::unbounded_channel::<String>();
214
215        // The request carries an `OnUpgrade` extension, as the http server sets.
216        let (pending_upgrade, on_upgrade) = pending();
217        let req = Request::new(Body::empty());
218        req.extensions().insert(on_upgrade);
219
220        // Responder echoes the request back (so `handle_upgrade` can find the
221        // `OnUpgrade`) and yields a response.
222        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        // Handler that always fails — previously this had to be `Infallible`.
231        let handler = service_fn(|_upgraded: Upgraded| async move {
232            Err::<(), BoxError>(BoxError::from_static_str("handler boom"))
233        });
234
235        // Fallthrough inner service (not reached: matcher is `true`).
236        let inner =
237            service_fn(
238                |_req: Request| async move { Ok::<_, Infallible>(Response::new(Body::empty())) },
239            );
240
241        // The handler keeps its own error type; its (raw) error is routed to
242        // the per-handler sink given here.
243        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        // Serving spawns the detached upgrade task (which awaits the upgrade).
255        let _resp = svc.serve(req).await.expect("upgrade match -> Ok(response)");
256
257        // Fulfill the pending upgrade so the handler runs and then fails.
258        let upgraded = Upgraded::new(ServiceInput::new(Builder::default().build()), Bytes::new());
259        pending_upgrade.fulfill(upgraded);
260
261        // The handler error must reach the sink (not be swallowed).
262        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}