rama_http/layer/har/recorder/
fs.rs1use super::Recorder;
2use crate::layer::har::spec;
3use rama_core::error::{BoxError, ErrorContext};
4use rama_core::extensions::{Extension, Extensions};
5use rama_core::telemetry::tracing;
6use rama_utils::{
7 fs::{CreatedFilePermissions, OpenOptions},
8 time::now_unix,
9};
10use std::io::Write;
11use std::ops::Deref;
12use std::path::PathBuf;
13use std::sync::Arc;
14use tokio::fs::File;
15use tokio::io::AsyncWriteExt;
16use tokio::sync::{mpsc, oneshot};
17use tokio::time::Instant;
18
19#[derive(Debug, Clone)]
22pub struct FileRecorder {
23 tx: mpsc::Sender<FileRecorderMessage>,
24}
25
26#[derive(Debug, Clone, Extension)]
27#[extension(tags(http))]
28pub struct HarFilePath(Arc<PathBuf>);
32
33impl AsRef<std::path::Path> for HarFilePath {
34 fn as_ref(&self) -> &std::path::Path {
35 self.0.as_ref()
36 }
37}
38
39impl Deref for HarFilePath {
40 type Target = std::path::Path;
41
42 fn deref(&self) -> &Self::Target {
43 self.0.as_ref()
44 }
45}
46
47#[derive(Debug)]
48enum FileRecorderMessage {
49 Record {
50 log: Box<spec::Log>,
51 ext: oneshot::Sender<Extensions>,
52 },
53 Stop,
54}
55
56#[derive(Debug)]
57struct FileRecorderTask {
58 rx: mpsc::Receiver<FileRecorderMessage>,
59
60 dir: PathBuf,
61 prefix: String,
62 start: Instant,
63 start_epoch: i64,
64}
65
66impl FileRecorderTask {
67 fn new(rx: mpsc::Receiver<FileRecorderMessage>, dir: PathBuf, prefix: String) -> Self {
68 Self {
69 rx,
70 dir,
71 prefix,
72 start: Instant::now(),
73 start_epoch: now_unix(),
74 }
75 }
76
77 async fn run(mut self) {
78 #[derive(Debug)]
79 struct Storage {
80 file: File,
81 path: PathBuf,
82 has_entries: bool,
83 }
84
85 impl Storage {
86 async fn try_new(path: PathBuf) -> Result<Self, BoxError> {
87 if let Some(parent) = path.parent() {
88 create_har_parent_dir(parent)
89 .await
90 .context("create HAR file parent dir")?;
91 }
92 let file = OpenOptions::new()
97 .write(true)
98 .create(true)
99 .truncate(true)
100 .created_file_permissions(CreatedFilePermissions::OwnerReadWrite)
101 .open(&path)
102 .await
103 .context("create HAR file")?;
104 Ok(Self {
105 file,
106 path,
107 has_entries: false,
108 })
109 }
110 }
111
112 async fn create_har_parent_dir(parent: &std::path::Path) -> std::io::Result<()> {
117 #[cfg(unix)]
118 {
119 use std::os::unix::fs::DirBuilderExt as _;
120 let mut b = std::fs::DirBuilder::new();
121 b.recursive(true).mode(0o700);
122 let parent = parent.to_owned();
123 tokio::task::spawn_blocking(move || b.create(&parent))
124 .await
125 .map_err(std::io::Error::other)??;
126 Ok(())
127 }
128 #[cfg(not(unix))]
129 {
130 tokio::fs::create_dir_all(parent).await
131 }
132 }
133
134 let mut storage: Option<Storage> = None;
135 let mut counter = 0;
136 let mut buf = Vec::new();
137
138 'msg_loop: while let Some(msg) = self.rx.recv().await {
139 match msg {
140 FileRecorderMessage::Record { log, ext } => {
141 let storage_ref = if let Some(sr) = storage.as_mut() {
142 sr
143 } else {
144 storage = Some(
145 match async {
146 let file_name = format!(
147 "{}_{}_{}_{}.har",
148 self.prefix,
149 self.start_epoch,
150 {
151 let i = counter;
152 counter += 1;
153 i
154 },
155 self.start.elapsed().as_secs()
156 );
157 create_har_parent_dir(&self.dir)
158 .await
159 .context("create HAR recording dir")?;
160 let path = rama_utils::fs::safe_path_in(&self.dir, file_name)
161 .await
162 .context("validate HAR file path")?;
163 Storage::try_new(path).await
164 }
165 .await
166 {
167 Err(err) => {
168 tracing::debug!(
169 "failed to create file for HAR recording: {err} (ignore log entry)"
170 );
171 continue 'msg_loop;
172 }
173 Ok(storage) => storage,
174 },
175 );
176 #[expect(
177 clippy::expect_used,
178 reason = "it was assigned in the previous assignment as Some"
179 )]
180 let storage_ref = storage
182 .as_mut()
183 .expect("storage to be some due to previous statement");
184
185 buf.clear();
186 let header = serde_json::json!({
187 "log": {
188 "version": log.version,
189 "creator": log.creator,
190 "browser": log.browser,
191 "comment": log.comment,
192 "pages": [], },
194 });
195 if let Err(err) = serde_json::to_writer(&mut buf, &header) {
196 tracing::debug!(
197 "failed to serialize initial json content for HAR log: {err} (drop file)"
198 );
199 storage = None;
200 continue 'msg_loop;
201 }
202 buf.truncate(buf.len() - 2); _ = write!(buf, ",\"entries\":["); if let Err(err) = storage_ref.file.write_all(&buf).await {
205 tracing::debug!(
206 "failed to write initial json content for HAR log: {err} (drop file)"
207 );
208 storage = None;
209 continue 'msg_loop;
210 }
211
212 storage_ref
213 };
214
215 if log.pages.map(|p| !p.is_empty()).unwrap_or_default() {
216 tracing::debug!(
217 "log contains pages which are not supported by the har recorder!"
218 );
219 }
220
221 for entry in log.entries.iter() {
222 tracing::trace!("har log file writer: write entry: {entry:?}");
223 buf.clear();
224 match serde_json::to_writer(&mut buf, entry) {
225 Ok(_) => {
226 if storage_ref.has_entries
227 && let Err(err) = storage_ref.file.write_u8(b',').await
228 {
229 tracing::debug!("failed to write entry separator: {err}");
230 #[expect(clippy::expect_used)]
231 finish_file(
232 storage
233 .take()
234 .expect("storage to exist as we have reference to it")
235 .file,
236 )
237 .await;
238 continue 'msg_loop;
239 } else if let Err(err) = storage_ref.file.write_all(&buf).await {
240 tracing::debug!("failed to write serialized entry: {err}");
241 #[expect(clippy::expect_used)]
242 finish_file(
243 storage
244 .take()
245 .expect("storage to exist as we have reference to it")
246 .file,
247 )
248 .await;
249 continue 'msg_loop;
250 } else {
251 storage_ref.has_entries = true;
252 }
253 }
254 Err(err) => {
255 tracing::debug!(
256 "failed entry ({entry:?}) due to json serialize error: {err}"
257 );
258 #[expect(clippy::expect_used)]
259 finish_file(
260 storage
261 .take()
262 .expect("storage to exist as we have reference to it")
263 .file,
264 )
265 .await;
266 continue 'msg_loop;
267 }
268 }
269 }
270
271 let extensions = Extensions::new();
272 extensions.insert(HarFilePath(storage_ref.path.clone().into()));
273 if ext.send(extensions).is_err() {
274 tracing::debug!(
275 "failed to send http extensions w/ har file path back to recorder callee"
276 );
277 }
278 }
279 FileRecorderMessage::Stop => {
280 if let Some(storage) = storage.take() {
281 tracing::trace!(
282 "FileRecorderMessage::Stop received: finish file {:?}",
283 storage.path
284 );
285 finish_file(storage.file).await;
286 } else {
287 tracing::debug!(
288 "FileRecorderMessage::Stop received while no session active: ignore"
289 );
290 }
291 }
292 }
293 }
294 if let Some(storage) = storage {
295 tracing::trace!(
296 "FileRecorder task exiting: file '{:?}' was still active: finish file",
297 storage.path
298 );
299 finish_file(storage.file).await;
300 }
301 }
302}
303
304async fn finish_file(mut file: File) {
305 if let Err(err) = file.write_all(b"]}}").await {
307 tracing::debug!("failed to write trailing characters for finished har file: {err}");
308 }
309}
310
311impl Default for FileRecorder {
312 fn default() -> Self {
313 Self::new(
314 std::env::temp_dir().join("rama").join("har_recordings"),
315 format!(
316 "rama_{}_recording",
317 rama_utils::info::VERSION.replace('.', "_")
318 ),
319 )
320 }
321}
322
323impl FileRecorder {
324 #[must_use]
329 pub fn new(dir: PathBuf, prefix: String) -> Self {
330 let (tx, rx) = mpsc::channel(match std::thread::available_parallelism() {
331 Ok(n) => n.get(),
332 Err(_) => 1,
333 });
334
335 let task = FileRecorderTask::new(rx, dir, prefix);
336 tokio::spawn(task.run());
337
338 Self { tx }
339 }
340}
341
342impl Recorder for FileRecorder {
343 async fn record(&self, log: spec::Log) -> Option<Extensions> {
344 let (tx, rx) = oneshot::channel();
345 if let Err(err) = self
346 .tx
347 .send(FileRecorderMessage::Record {
348 log: Box::new(log),
349 ext: tx,
350 })
351 .await
352 {
353 tracing::debug!("FileRecorder: failed to send log for recording to task: {err}");
354 }
355 rx.await
356 .inspect_err(|err| {
357 tracing::debug!("file recorder: record oneshot reply await error: {err}");
358 })
359 .ok()
360 }
361
362 async fn stop_record(&self) {
363 if let Err(err) = self.tx.send(FileRecorderMessage::Stop).await {
364 tracing::debug!("FileRecorder: failed to send stop record msg to task: {err}");
365 }
366 }
367}