rama_http/layer/har/
service.rs1use 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 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, request: HarRequest,
60 }
61
62 let (request, maybe_entry_start_info) = if self.toggle.status().await {
63 let start_time = Timestamp::now();
64 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 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 server_ip_address: None,
156 connection: None, 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}