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 tokio::time::sleep(Duration::from_millis(50)).await;
787
788 logger.flush().await.unwrap();
789
790 tokio::time::sleep(Duration::from_millis(50)).await;
792
793 logger.shutdown().await;
794
795 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}