Skip to main content

rama_http/layer/har/
service.rs

1use crate::body::util::BodyExt;
2use crate::layer::har::recorder::Recorder;
3use crate::layer::har::spec::{
4    Cache, Entry, Log as HarLog, Request as HarRequest, Response as HarResponse, Timings,
5};
6use crate::layer::har::toggle::Toggle;
7use crate::{Body, Request, Response, StreamingBody};
8
9use jiff::Timestamp;
10
11use rama_core::error::{BoxError, ErrorContext as _};
12use rama_core::extensions::ExtensionsRef;
13use rama_core::telemetry::tracing;
14use rama_core::{Service, bytes::Bytes};
15use tokio::time::Instant;
16
17pub struct HARExportService<R, S, T> {
18    pub(super) recorder: R,
19    pub(super) service: S,
20    pub(super) toggle: T,
21
22    pub(super) preserve_sensitive: bool,
23}
24
25impl<R, S, T> HARExportService<R, S, T> {
26    pub fn recorder(&self) -> &R {
27        &self.recorder
28    }
29
30    pub fn toggle(&self) -> &T {
31        &self.toggle
32    }
33
34    rama_utils::macros::generate_set_and_with! {
35        /// Sets whether to preserve sensitive headers (false by default).
36        pub fn preserve_sensitive(mut self) -> Self {
37            self.preserve_sensitive = true;
38            self
39        }
40    }
41}
42
43impl<R, S, W, ReqBody, ResBody> Service<Request<ReqBody>> for HARExportService<R, S, W>
44where
45    R: Recorder,
46    S: Service<Request, Output = Response<ResBody>>,
47    S::Error: Into<BoxError> + Send + Sync + 'static,
48    W: Toggle,
49    ReqBody: StreamingBody<Data = Bytes, Error: Into<BoxError>> + Send + Sync + 'static,
50    ResBody: StreamingBody<Data = Bytes, Error: Into<BoxError>> + Send + Sync + 'static,
51{
52    type Output = Response;
53    type Error = BoxError;
54
55    async fn serve(&self, req: Request<ReqBody>) -> Result<Self::Output, Self::Error> {
56        struct EntryStartInfo {
57            start_time: Timestamp,
58            begin: Instant, // TODO: replace with total time
59            request: HarRequest,
60        }
61
62        let (request, maybe_entry_start_info) = if self.toggle.status().await {
63            let start_time = Timestamp::now();
64            // need to collect it first as bodies are (potentially) streaming
65            let (req_parts, req_body) = req.into_parts();
66            let req_body_bytes = req_body
67                .collect()
68                .await
69                .context("collect request body for HAR recording and inner svc")?
70                .to_bytes();
71
72            let har_req_result = HarRequest::from_http_request_parts(
73                &req_parts,
74                &req_body_bytes,
75                self.preserve_sensitive,
76            );
77            let request = Request::from_parts(req_parts, Body::from(req_body_bytes));
78
79            match har_req_result {
80                Err(err) => {
81                    tracing::debug!(
82                        "failed to create HAR request from incoming HTTP Request: {err}"
83                    );
84                    (request, None)
85                }
86                Ok(har_request) => {
87                    let info = EntryStartInfo {
88                        start_time,
89                        begin: Instant::now(),
90                        request: har_request,
91                    };
92                    (request, Some(info))
93                }
94            }
95        } else {
96            self.recorder.stop_record().await;
97            (req.map(Body::new), None)
98        };
99
100        let result = self.service.serve(request).await;
101
102        if let Some(entry_start_info) = maybe_entry_start_info {
103            let (result, response) = match result {
104                Ok(resp) => {
105                    let (resp_parts, resp_body) = resp.into_parts();
106                    let resp_body_bytes = resp_body
107                        .collect()
108                        .await
109                        .context("collect response body for HAR recording and return value")?
110                        .to_bytes();
111
112                    let maybe_response = match HarResponse::from_http_response_parts(
113                        &resp_parts,
114                        &resp_body_bytes,
115                        self.preserve_sensitive,
116                    ) {
117                        Err(err) => {
118                            tracing::debug!(
119                                "failed to create HAR response from returned HTTP Response: {err}"
120                            );
121                            None
122                        }
123                        Ok(resp) => Some(resp),
124                    };
125
126                    let result = Ok(Response::from_parts(
127                        resp_parts,
128                        Body::from(resp_body_bytes),
129                    ));
130
131                    (result, maybe_response)
132                }
133                Err(err) => (Err(err.into()), None),
134            };
135
136            // TODO: populate these in future
137            let timings = Timings::default();
138            let cache = Cache::default();
139
140            let entry = Entry {
141                page_ref: None,
142                started_date_time: entry_start_info.start_time,
143                time: entry_start_info
144                    .begin
145                    .elapsed()
146                    .as_millis()
147                    .min(i64::MAX as u128) as i64,
148                request: entry_start_info.request,
149                response,
150                cache,
151                timings,
152                // TODO: when used as server middleware it is SocketInfo local addr,
153                //       but when used via client middleware it is supposed to be the resolved address,
154                //       which I am not sure is already exposed (TODO^2)
155                server_ip_address: None,
156                connection: None, // TODO
157                comment: None,
158            };
159
160            let log_line = HarLog {
161                entries: vec![entry],
162                ..Default::default()
163            };
164
165            let maybe_resp_extensions = self.recorder.record(log_line).await;
166
167            let result = match (result, maybe_resp_extensions) {
168                (Ok(resp), Some(resp_extensions)) => {
169                    tracing::trace!("extend (ok) response with HAR recorder extensions");
170                    resp.extensions().extend(&resp_extensions);
171                    Ok(resp)
172                }
173                (result, _) => result,
174            };
175
176            return result;
177        }
178
179        match result {
180            Ok(response) => Ok(response.map(Body::new)),
181            Err(err) => Err(err.into()),
182        }
183    }
184}