Skip to main content

kftray_http_logs/
logger.rs

1use std::path::PathBuf;
2use std::sync::Arc;
3
4use anyhow::{
5    Context,
6    Result,
7};
8use bytes::{
9    BufMut,
10    Bytes,
11    BytesMut,
12};
13use chrono::{
14    DateTime,
15    Utc,
16};
17use dashmap::DashMap;
18use lazy_static::lazy_static;
19use tokio::fs::{
20    File,
21    OpenOptions,
22};
23use tokio::io::{
24    AsyncWriteExt,
25    BufWriter,
26};
27use tokio::sync::RwLock;
28use tokio::sync::mpsc::{
29    self,
30    Sender,
31};
32use tokio::time::Duration;
33use tracing::{
34    debug,
35    error,
36    trace,
37};
38use uuid::Uuid;
39
40use crate::config::LogConfig;
41use crate::formatter::MessageFormatter;
42use crate::message::LogMessage;
43
44#[derive(Debug, Clone)]
45pub struct TraceInfo {
46    pub trace_id: String,
47    pub timestamp: DateTime<Utc>,
48}
49
50pub fn calculate_time_diff(start: DateTime<Utc>, end: DateTime<Utc>) -> i64 {
51    (end - start).num_milliseconds()
52}
53
54lazy_static! {
55    static ref BUFFER_POOL: Arc<tokio::sync::Mutex<Vec<BytesMut>>> =
56        Arc::new(tokio::sync::Mutex::new(Vec::with_capacity(32)));
57}
58
59const CHANNEL_CAPACITY: usize = 256;
60const BATCH_SIZE_THRESHOLD: usize = 10;
61const FLUSH_INTERVAL_MS: u64 = 100;
62const TRACE_CLEANUP_INTERVAL_SECS: u64 = 5;
63const TRACE_EXPIRY_SECS: i64 = 1800;
64
65type TraceMap = Arc<DashMap<String, TraceInfo>>;
66
67#[derive(Clone, Debug)]
68pub struct HttpLogger {
69    log_sender: Sender<LogMessage>,
70    trace_map: TraceMap,
71    shutdown: Arc<tokio::sync::watch::Sender<()>>,
72    #[allow(dead_code)]
73    config: LogConfig,
74    #[allow(dead_code)]
75    writer_task: Arc<tokio::sync::Mutex<Option<tokio::task::JoinHandle<()>>>>,
76    #[allow(dead_code)]
77    cleanup_task: Arc<tokio::sync::Mutex<Option<tokio::task::JoinHandle<()>>>>,
78}
79
80impl HttpLogger {
81    pub async fn new(log_config: LogConfig, log_file_path: PathBuf) -> Result<Self> {
82        let (log_sender, mut log_receiver) = mpsc::channel::<LogMessage>(CHANNEL_CAPACITY);
83
84        let log_file = Arc::new(RwLock::new(BufWriter::with_capacity(
85            64 * 1024,
86            OpenOptions::new()
87                .append(true)
88                .create(true)
89                .open(&log_file_path)
90                .await
91                .context("Failed to open log file")?,
92        )));
93
94        let trace_map: TraceMap = Arc::new(DashMap::with_capacity(1024));
95        let (shutdown_tx, mut shutdown_rx) = tokio::sync::watch::channel(());
96        let mut shutdown_rx_writer = shutdown_rx.clone();
97
98        let writer_task = tokio::spawn({
99            let log_file = log_file.clone();
100            async move {
101                let mut flush_interval =
102                    tokio::time::interval(Duration::from_millis(FLUSH_INTERVAL_MS));
103                let mut message_batch = Vec::with_capacity(BATCH_SIZE_THRESHOLD * 2);
104                let mut last_flush = Utc::now();
105
106                loop {
107                    tokio::select! {
108                        Some(log_message) = log_receiver.recv() => {
109                            if log_message.is_flush_trigger() {
110                                if let Ok(mut file) = log_file.try_write() {
111                                    let _ = file.flush().await;
112                                }
113                                continue;
114                            }
115
116                            let is_response = log_message.is_response();
117                            message_batch.push(log_message);
118
119                            let now = Utc::now();
120                            let batch_too_old = now.signed_duration_since(last_flush).num_milliseconds() > FLUSH_INTERVAL_MS as i64;
121                            let force_write = message_batch.len() >= BATCH_SIZE_THRESHOLD || is_response || batch_too_old;
122
123                            if force_write {
124                                let batch_size = message_batch.len();
125                                debug!("Writing log batch of {} messages (contains response: {})",
126                                      batch_size, is_response);
127
128                                match Self::write_log_batch(&log_file, &message_batch).await { Err(e) => {
129                                    error!("Failed to write log batch: {:?}", e);
130                                } _ => {
131                                    debug!("Successfully wrote log batch");
132
133                                    if is_response
134                                        && let Ok(mut file) = log_file.try_write()
135                                        && let Err(e) = file.get_mut().sync_data().await {
136                                                error!("Failed to sync response log to disk: {:?}", e);
137                                            }
138                                }}
139                                message_batch.clear();
140                                last_flush = Utc::now();
141                            }
142                        }
143                        _ = flush_interval.tick() => {
144                            let now = Utc::now();
145                            if !message_batch.is_empty() {
146                                debug!("Timer flush: writing batch of {} messages after {}ms",
147                                      message_batch.len(),
148                                      now.signed_duration_since(last_flush).num_milliseconds());
149
150                                if let Err(e) = Self::write_log_batch(&log_file, &message_batch).await {
151                                    error!("Failed to write log batch: {:?}", e);
152                                }
153                                message_batch.clear();
154                                last_flush = now;
155                            }
156
157                            if let Ok(mut file) = log_file.try_write() {
158                                if let Err(e) = file.flush().await {
159                                    error!("Failed to flush log file in timer: {:?}", e);
160                                } else {
161                                    trace!("Successfully flushed log file during periodic tick");
162                                }
163                            }
164                        }
165                        _ = shutdown_rx_writer.changed() => {
166                            debug!("Shutting down log writer - processing final messages");
167
168                            if !message_batch.is_empty() {
169                                debug!("Writing final batch of {} messages during shutdown", message_batch.len());
170
171                                if let Err(e) = Self::write_log_batch(&log_file, &message_batch).await {
172                                    error!("Failed to write final log batch: {:?}", e);
173
174                                    for (i, msg) in message_batch.iter().enumerate() {
175                                        debug!("Attempting to write individual message {} during shutdown", i+1);
176                                        if let Err(e) = Self::write_single_log(&log_file, msg).await {
177                                            error!("Failed to write message {}: {:?}", i+1, e);
178                                        }
179                                    }
180                                }
181                            }
182
183                            debug!("Performing final sync to ensure data durability");
184                            let mut file = log_file.write().await;
185                            if let Err(e) = file.flush().await {
186                                error!("Failed to flush log file during shutdown: {:?}", e);
187                            } else if let Err(e) = file.get_mut().sync_all().await {
188                                error!("Failed to sync log file during shutdown: {:?}", e);
189                            } else {
190                                debug!("Successfully flushed and synced log file during shutdown");
191                            }
192
193                            debug!("Log writer shutdown complete");
194                            break;
195                        }
196                    }
197                }
198            }
199        });
200
201        let cleanup_task = tokio::spawn({
202            let trace_map = trace_map.clone();
203            async move {
204                let mut interval =
205                    tokio::time::interval(Duration::from_secs(TRACE_CLEANUP_INTERVAL_SECS));
206                loop {
207                    tokio::select! {
208                        _ = interval.tick() => {
209                            let now = Utc::now();
210                            trace_map.retain(|_, trace_info| {
211                                now.signed_duration_since(trace_info.timestamp).num_seconds() < TRACE_EXPIRY_SECS
212                            });
213                        }
214                        _ = shutdown_rx.changed() => {
215                            debug!("Shutting down cleanup task");
216                            break;
217                        }
218                    }
219                }
220            }
221        });
222
223        let writer_task_handle = Arc::new(tokio::sync::Mutex::new(Some(writer_task)));
224        let cleanup_task_handle = Arc::new(tokio::sync::Mutex::new(Some(cleanup_task)));
225
226        Ok(Self {
227            log_sender,
228            trace_map,
229            shutdown: Arc::new(shutdown_tx),
230            config: log_config,
231            writer_task: writer_task_handle,
232            cleanup_task: cleanup_task_handle,
233        })
234    }
235
236    pub async fn for_config(config_id: i64, local_port: u16) -> Result<Self> {
237        let log_config = LogConfig::new(LogConfig::default_log_directory()?);
238        let log_path = log_config
239            .create_log_file_path(config_id, local_port)
240            .await?;
241        Self::new(log_config, log_path).await
242    }
243
244    pub async fn log_request(&self, buffer: Bytes) -> String {
245        let request_id = Uuid::new_v4().to_string();
246        let timestamp = Utc::now();
247        let trace_id = request_id.clone();
248
249        if let Err(e) = self
250            .send_request_log(buffer, trace_id.clone(), timestamp)
251            .await
252        {
253            error!("Failed to log request: {:?}", e);
254        }
255
256        self.trace_map.insert(
257            request_id.clone(),
258            TraceInfo {
259                trace_id,
260                timestamp,
261            },
262        );
263
264        request_id
265    }
266
267    pub async fn log_response(&self, buffer: Bytes, request_id: String) {
268        let timestamp = Utc::now();
269        let is_preformatted = buffer.len() > 5 && &buffer[0..5] == b"HTTP/";
270
271        match self.trace_map.get(&request_id) {
272            Some(trace_info) => {
273                let took_ms = calculate_time_diff(trace_info.timestamp, timestamp);
274                self.send_response_log_internal(
275                    buffer,
276                    request_id,
277                    timestamp,
278                    took_ms,
279                    is_preformatted,
280                )
281                .await;
282            }
283            _ => {
284                debug!("No trace info found for request ID: {}", request_id);
285                self.send_response_log_internal(buffer, request_id, timestamp, 0, is_preformatted)
286                    .await;
287            }
288        }
289    }
290
291    async fn send_response_log_internal(
292        &self, buffer: Bytes, request_id: String, timestamp: DateTime<Utc>, took_ms: i64,
293        is_preformatted: bool,
294    ) {
295        let result = if is_preformatted {
296            self.send_preformatted_response_log(buffer, request_id.clone(), timestamp, took_ms)
297                .await
298        } else {
299            self.send_response_log(buffer, request_id.clone(), timestamp, took_ms)
300                .await
301        };
302
303        if let Err(e) = result {
304            error!("Failed to send response log: {:?}", e);
305        }
306    }
307
308    async fn send_request_log(
309        &self, buffer: Bytes, trace_id: String, timestamp: DateTime<Utc>,
310    ) -> Result<()> {
311        debug!("Formatting request log for trace ID: {}", trace_id);
312
313        let log_entry = match MessageFormatter::format_request(&buffer, &trace_id, timestamp).await
314        {
315            Ok(entry) => entry,
316            Err(e) => {
317                debug!("Failed to parse HTTP request for logging: {:?}", e);
318                return Ok(());
319            }
320        };
321
322        match self.log_sender.send(log_entry).await {
323            Ok(_) => debug!("Successfully sent request log message to channel"),
324            Err(e) => error!("Failed to send log message: {:?}", e),
325        }
326
327        Ok(())
328    }
329
330    async fn send_response_log(
331        &self, buffer: Bytes, trace_id: String, timestamp: DateTime<Utc>, took_ms: i64,
332    ) -> Result<()> {
333        debug!(
334            "Formatting response log for trace ID: {} (size: {}B)",
335            trace_id,
336            buffer.len()
337        );
338
339        let is_valid_http =
340            buffer.len() > 16 && matches!(buffer.get(..5), Some(prefix) if prefix == b"HTTP/");
341
342        if !is_valid_http {
343            debug!("Response doesn't appear to be a valid HTTP response, but will log anyway");
344        }
345
346        let log_entry =
347            match MessageFormatter::format_response(&buffer, &trace_id, timestamp, took_ms).await {
348                Ok(entry) => entry,
349                Err(e) => {
350                    debug!("Failed to parse HTTP response for logging: {:?}", e);
351                    return Ok(());
352                }
353            };
354
355        let message_size = log_entry.size();
356        match self.log_sender.send(log_entry).await {
357            Ok(_) => debug!(
358                "Successfully sent response log message to channel (size: {}B)",
359                message_size
360            ),
361            Err(e) => {
362                error!(
363                    "Failed to send response log message (size: {}B): {:?}",
364                    message_size, e
365                );
366                return Err(anyhow::anyhow!(
367                    "Failed to send response log message: {:?}",
368                    e
369                ));
370            }
371        }
372
373        Ok(())
374    }
375
376    async fn send_preformatted_response_log(
377        &self, buffer: Bytes, trace_id: String, timestamp: DateTime<Utc>, took_ms: i64,
378    ) -> Result<()> {
379        let message = LogMessage::new_preformatted_response(trace_id, timestamp, took_ms, buffer);
380
381        if let Err(e) = self.log_sender.send(message).await {
382            error!("Failed to send preformatted response log message: {:?}", e);
383            return Err(anyhow::anyhow!("Failed to send log message"));
384        }
385
386        Ok(())
387    }
388
389    pub async fn flush(&self) -> Result<()> {
390        let trigger = LogMessage::TriggerFlush;
391        self.log_sender
392            .send(trigger)
393            .await
394            .map_err(|_| anyhow::anyhow!("Failed to send flush trigger"))
395    }
396
397    async fn write_log_batch(
398        log_file: &Arc<RwLock<BufWriter<File>>>, messages: &[LogMessage],
399    ) -> Result<()> {
400        if messages.is_empty() {
401            return Ok(());
402        }
403
404        let mut total_size = 0;
405        let mut response_count = 0;
406
407        for message in messages {
408            total_size += message.as_bytes().len();
409            if message.is_response() {
410                response_count += 1;
411            }
412        }
413
414        let mut combined_buffer = BytesMut::with_capacity(total_size);
415
416        if response_count > 0 {
417            debug!("Processing {} response messages in batch", response_count);
418
419            for message in messages.iter() {
420                if message.is_response() {
421                    let bytes = message.as_bytes();
422                    combined_buffer.put_slice(bytes);
423                    trace!(
424                        "Added response message to write buffer: {} bytes",
425                        bytes.len()
426                    );
427                }
428            }
429        }
430
431        for message in messages.iter() {
432            if !message.is_response() {
433                combined_buffer.put_slice(message.as_bytes());
434            }
435        }
436
437        let mut log_file = log_file.write().await;
438        debug!(
439            "Acquired write lock for log file batch of {} messages (buffer size: {}B)",
440            messages.len(),
441            combined_buffer.len()
442        );
443
444        log_file
445            .write_all(&combined_buffer)
446            .await
447            .context("Failed to write log entries to file")?;
448
449        log_file
450            .flush()
451            .await
452            .context("Failed to flush log entries to file")?;
453
454        if response_count > 0 {
455            debug!("Syncing {} response messages to disk", response_count);
456
457            if let Err(e) = log_file.get_mut().sync_data().await {
458                error!("Failed to sync log file data to disk: {:?}", e);
459            } else {
460                debug!("Successfully synced log file with responses to disk");
461            }
462        }
463
464        debug!(
465            "Successfully wrote and flushed batch of {} messages (responses: {}, total bytes: {})",
466            messages.len(),
467            response_count,
468            combined_buffer.len()
469        );
470
471        Ok(())
472    }
473
474    async fn write_single_log(
475        log_file: &Arc<RwLock<BufWriter<File>>>, message: &LogMessage,
476    ) -> Result<()> {
477        let mut log_file = log_file.write().await;
478        trace!(
479            "Acquired write lock for single {} message",
480            message.message_type()
481        );
482
483        log_file
484            .write_all(message.as_bytes())
485            .await
486            .context("Failed to write single log entry to file")?;
487
488        log_file
489            .flush()
490            .await
491            .context("Failed to flush single log entry to file")?;
492
493        trace!(
494            "Wrote and flushed single {} message",
495            message.message_type()
496        );
497
498        Ok(())
499    }
500
501    pub async fn shutdown(&self) {
502        debug!("Initiating HTTP logger shutdown sequence");
503
504        let _ = self.flush().await;
505
506        let shutdown_signal = self.shutdown.clone();
507
508        let _ = shutdown_signal.send(());
509        debug!("Sent shutdown signal to logger tasks");
510
511        let timeout = Duration::from_secs(1);
512
513        let writer_handle = {
514            let mut guard = self.writer_task.lock().await;
515            guard.take()
516        };
517
518        if let Some(handle) = writer_handle {
519            debug!("Awaiting writer task completion");
520            match tokio::time::timeout(timeout, handle).await {
521                Ok(result) => {
522                    if let Err(e) = result {
523                        error!("Error awaiting writer task: {:?}", e);
524                    }
525                }
526                Err(_) => error!("Writer task shutdown timed out"),
527            }
528        }
529
530        let cleanup_handle = {
531            let mut guard = self.cleanup_task.lock().await;
532            guard.take()
533        };
534
535        if let Some(handle) = cleanup_handle {
536            debug!("Awaiting cleanup task completion");
537            match tokio::time::timeout(timeout, handle).await {
538                Ok(result) => {
539                    if let Err(e) = result {
540                        error!("Error awaiting cleanup task: {:?}", e);
541                    }
542                }
543                Err(_) => error!("Cleanup task shutdown timed out"),
544            }
545        }
546
547        debug!("HTTP logger shutdown sequence completed");
548    }
549}
550
551#[cfg(test)]
552mod tests {
553    #[allow(unused_imports)]
554    use std::io::Write;
555
556    use chrono::Utc;
557    use mockall::mock;
558    use tempfile::tempdir;
559    use tokio::io::AsyncWriteExt;
560
561    use super::*;
562
563    mock! {
564        pub MockFile {}
565    }
566
567    #[test]
568    fn test_trace_map_operations() {
569        let trace_map: Arc<DashMap<String, TraceInfo>> = Arc::new(DashMap::with_capacity(10));
570
571        let old_time = Utc::now() - chrono::Duration::seconds(TRACE_EXPIRY_SECS + 10);
572
573        for i in 0..5 {
574            let trace_id = format!("old-trace-{i}");
575            trace_map.insert(
576                trace_id.clone(),
577                TraceInfo {
578                    trace_id,
579                    timestamp: old_time,
580                },
581            );
582        }
583
584        for i in 0..5 {
585            let trace_id = format!("recent-trace-{i}");
586            trace_map.insert(
587                trace_id.clone(),
588                TraceInfo {
589                    trace_id,
590                    timestamp: Utc::now(),
591                },
592            );
593        }
594
595        assert_eq!(trace_map.len(), 10);
596
597        let now = Utc::now();
598        trace_map.retain(|_, info| {
599            let age_secs = (now - info.timestamp).num_seconds();
600            age_secs <= TRACE_EXPIRY_SECS
601        });
602
603        assert_eq!(trace_map.len(), 5);
604
605        for i in 0..5 {
606            let old_id = format!("old-trace-{i}");
607            let recent_id = format!("recent-trace-{i}");
608
609            assert!(!trace_map.contains_key(&old_id));
610            assert!(trace_map.contains_key(&recent_id));
611        }
612    }
613
614    #[test]
615    fn test_trace_info() {
616        let trace_id = "test-trace-123".to_string();
617        let timestamp = Utc::now();
618
619        let info = TraceInfo {
620            trace_id: trace_id.clone(),
621            timestamp,
622        };
623
624        assert_eq!(info.trace_id, trace_id);
625        assert_eq!(info.timestamp, timestamp);
626    }
627
628    #[test]
629    fn test_calculate_time_diff() {
630        let start = Utc::now();
631        let end = start + chrono::Duration::milliseconds(500);
632
633        let diff = calculate_time_diff(start, end);
634        assert_eq!(diff, 500);
635    }
636
637    #[test]
638    fn test_log_message_types() {
639        let req_content = "REQUEST: GET /test HTTP/1.1";
640        let resp_content = "RESPONSE: HTTP/1.1 200 OK";
641        let preformatted = "PREFORMATTED: Some special response";
642
643        let req_msg = LogMessage::Request(req_content.to_string());
644        let resp_msg = LogMessage::Response(resp_content.to_string());
645        let preformatted_msg = LogMessage::PreformattedResponse(preformatted.to_string());
646        let flush_msg = LogMessage::TriggerFlush;
647
648        assert_eq!(req_msg.message_type(), "Request");
649        assert_eq!(resp_msg.message_type(), "Response");
650        assert_eq!(preformatted_msg.message_type(), "PreformattedResponse");
651        assert_eq!(flush_msg.message_type(), "TriggerFlush");
652
653        assert!(!req_msg.is_response());
654        assert!(resp_msg.is_response());
655        assert!(preformatted_msg.is_response());
656        assert!(!flush_msg.is_response());
657
658        assert!(!req_msg.is_flush_trigger());
659        assert!(!resp_msg.is_flush_trigger());
660        assert!(!preformatted_msg.is_flush_trigger());
661        assert!(flush_msg.is_flush_trigger());
662
663        assert_eq!(req_msg.size(), req_content.len());
664        assert_eq!(resp_msg.size(), resp_content.len());
665        assert_eq!(preformatted_msg.size(), preformatted.len());
666        assert_eq!(flush_msg.size(), 0);
667    }
668
669    #[tokio::test]
670    async fn test_sender_receiver_pattern() {
671        let (tx, mut rx) = mpsc::channel::<LogMessage>(10);
672
673        let test_request = LogMessage::Request("Test request".to_string());
674        let test_response = LogMessage::Response("Test response".to_string());
675
676        tx.send(test_request.clone()).await.unwrap();
677        tx.send(test_response.clone()).await.unwrap();
678
679        let received_request = rx.recv().await.unwrap();
680        let received_response = rx.recv().await.unwrap();
681
682        if let LogMessage::Request(content) = received_request {
683            assert_eq!(content, "Test request");
684        } else {
685            panic!("Expected LogMessage::Request");
686        }
687
688        if let LogMessage::Response(content) = received_response {
689            assert_eq!(content, "Test response");
690        } else {
691            panic!("Expected LogMessage::Response");
692        }
693    }
694
695    #[test]
696    fn test_create_unique_trace_id() {
697        let mut trace_ids = Vec::with_capacity(100);
698
699        for _ in 0..100 {
700            let id = Uuid::new_v4().to_string();
701            trace_ids.push(id);
702        }
703
704        use std::collections::HashSet;
705        let trace_id_set: HashSet<_> = trace_ids.into_iter().collect();
706
707        assert_eq!(trace_id_set.len(), 100);
708    }
709
710    #[tokio::test]
711    async fn test_http_logger_creation() {
712        let temp_dir = tempdir().unwrap();
713        let file_path = temp_dir.path().join("test_log.txt");
714
715        let config = LogConfig::new(temp_dir.path().to_path_buf());
716        let logger = HttpLogger::new(config, file_path.clone()).await.unwrap();
717
718        assert!(logger.log_sender.capacity() >= CHANNEL_CAPACITY);
719        assert!(logger.trace_map.capacity() >= 1024);
720
721        logger.shutdown().await;
722    }
723
724    #[tokio::test]
725    async fn test_log_write_batch() {
726        let temp_dir = tempdir().unwrap();
727        let file_path = temp_dir.path().join("batch_test.log");
728
729        let file = File::create(&file_path).await.unwrap();
730        let buf_writer = BufWriter::new(file);
731
732        let log_file = Arc::new(RwLock::new(buf_writer));
733
734        let messages = vec![
735            LogMessage::Request("Test request 1".to_string()),
736            LogMessage::Request("Test request 2".to_string()),
737            LogMessage::Response("Test response".to_string()),
738        ];
739
740        HttpLogger::write_log_batch(&log_file, &messages)
741            .await
742            .unwrap();
743
744        {
745            let mut writer = log_file.write().await;
746            writer.flush().await.unwrap();
747        }
748
749        let contents = tokio::fs::read_to_string(&file_path).await.unwrap();
750
751        assert!(contents.contains("Test request 1"));
752        assert!(contents.contains("Test request 2"));
753        assert!(contents.contains("Test response"));
754
755        let response_pos = contents.find("Test response").unwrap();
756        let request1_pos = contents.find("Test request 1").unwrap();
757
758        assert!(
759            response_pos < request1_pos,
760            "Responses should appear before requests in the log file"
761        );
762    }
763
764    #[tokio::test]
765    async fn test_request_response_logging() {
766        let temp_dir = tempdir().unwrap();
767        let file_path = temp_dir.path().join("req_resp_test.log");
768
769        let config = LogConfig::new(temp_dir.path().to_path_buf());
770        let logger = HttpLogger::new(config, file_path.clone()).await.unwrap();
771
772        let request_data = Bytes::from_static(b"GET /test HTTP/1.1\r\nHost: example.com\r\n\r\n");
773        let response_data =
774            Bytes::from_static(b"HTTP/1.1 200 OK\r\nContent-Type: text/plain\r\n\r\nSuccess");
775
776        let request_id = logger.log_request(request_data).await;
777
778        assert!(!request_id.is_empty());
779        assert!(logger.trace_map.contains_key(&request_id));
780
781        tokio::time::sleep(Duration::from_millis(10)).await;
782
783        logger.log_response(response_data, request_id.clone()).await;
784
785        // Add a longer delay to ensure the response is written to disk
786        tokio::time::sleep(Duration::from_millis(50)).await;
787
788        logger.flush().await.unwrap();
789
790        // Add another small delay after flushing
791        tokio::time::sleep(Duration::from_millis(50)).await;
792
793        logger.shutdown().await;
794
795        // Add delay after shutdown
796        tokio::time::sleep(Duration::from_millis(50)).await;
797
798        let contents = tokio::fs::read_to_string(&file_path).await.unwrap();
799
800        println!("Log file contents: {contents}");
801
802        assert!(contents.contains("GET /test") || contents.contains("/test"));
803        assert!(contents.contains("200") || contents.contains("OK"));
804
805        let trace_info = logger.trace_map.get(&request_id).unwrap();
806        assert_eq!(trace_info.trace_id, request_id);
807    }
808
809    #[tokio::test]
810    async fn test_write_single_log() {
811        let temp_dir = tempdir().unwrap();
812        let file_path = temp_dir.path().join("single_log.txt");
813
814        let file = File::create(&file_path).await.unwrap();
815        let buf_writer = BufWriter::new(file);
816
817        let log_file = Arc::new(RwLock::new(buf_writer));
818
819        let message = LogMessage::Request("Single log test message".to_string());
820
821        HttpLogger::write_single_log(&log_file, &message)
822            .await
823            .unwrap();
824
825        {
826            let mut writer = log_file.write().await;
827            writer.flush().await.unwrap();
828        }
829
830        let contents = tokio::fs::read_to_string(&file_path).await.unwrap();
831        assert!(contents.contains("Single log test message"));
832    }
833}