1#![allow(unused)]
2
3use async_trait::async_trait;
4use chrono::{SecondsFormat, Utc};
5use once_cell::sync::Lazy;
6use serde::{Deserialize, Serialize};
7use std::env;
8use std::path::{Path, PathBuf};
9use std::sync::atomic::{AtomicBool, Ordering};
10use std::sync::{Arc, Mutex, RwLock};
11use tokio::sync::mpsc::Sender;
12use tracing::field::{Field, Visit};
13use tracing::{Event, Level, Subscriber};
14use tracing_appender::non_blocking::WorkerGuard;
15use tracing_appender::rolling::Rotation;
16use tracing_log::LogTracer;
17use tracing_subscriber::layer::{Context, Layer, SubscriberExt};
18use tracing_subscriber::{EnvFilter, util::SubscriberInitExt};
19use uuid::Uuid;
20
21use crate::utils::storage::{BlobStorage, Offloadable};
22
23#[derive(Serialize, Deserialize)]
24pub struct LogModel {
25 pub task_id: String,
26 #[serde(skip_serializing_if = "Option::is_none")]
27 pub request_id: Option<Uuid>,
28 pub status: String,
29 pub level: String,
30 pub message: String,
31 pub timestamp: String,
32 #[serde(skip_serializing_if = "Option::is_none")]
33 pub traceback: Option<String>,
34}
35
36#[async_trait]
37impl Offloadable for LogModel {
38 fn should_offload(&self, _threshold: usize) -> bool {
39 false
40 }
41 async fn offload(&mut self, _storage: &Arc<dyn BlobStorage>) -> crate::errors::Result<()> {
42 Ok(())
43 }
44 async fn reload(&mut self, _storage: &Arc<dyn BlobStorage>) -> crate::errors::Result<()> {
45 Ok(())
46 }
47}
48
49impl crate::utils::priority::Prioritizable for LogModel {
50 fn get_priority(&self) -> crate::utils::priority::Priority {
51 match self.level.to_lowercase().as_str() {
52 "error" | "fatal" => crate::utils::priority::Priority::High,
53 _ => crate::utils::priority::Priority::Low,
54 }
55 }
56}
57
58#[derive(Debug)]
59pub enum LogError {
60 Io(std::io::Error),
61 Send(String),
62 Init(tracing_appender::rolling::InitError),
63}
64
65impl From<std::io::Error> for LogError {
66 fn from(err: std::io::Error) -> Self {
67 Self::Io(err)
68 }
69}
70
71impl From<tracing_appender::rolling::InitError> for LogError {
72 fn from(err: tracing_appender::rolling::InitError) -> Self {
73 Self::Init(err)
74 }
75}
76
77impl std::fmt::Display for LogError {
78 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
79 match self {
80 LogError::Io(err) => write!(f, "{err}"),
81 LogError::Send(msg) => write!(f, "{msg}"),
82 LogError::Init(err) => write!(f, "{err}"),
83 }
84 }
85}
86
87impl std::error::Error for LogError {}
88
89#[derive(Debug, Clone, Serialize)]
90pub struct LogRecord {
91 pub time: String,
92 #[serde(skip)]
93 pub level: Level,
94 #[serde(rename = "level")]
95 pub level_name: String,
96 pub module: String,
97 pub message: String,
98 #[serde(skip_serializing_if = "Option::is_none")]
99 pub status: Option<String>,
100 #[serde(skip_serializing_if = "Option::is_none")]
101 pub event_type: Option<String>,
102 #[serde(skip_serializing_if = "Option::is_none")]
103 pub phase: Option<String>,
104 #[serde(skip_serializing_if = "Option::is_none")]
105 pub error_kind: Option<String>,
106 #[serde(skip_serializing_if = "Option::is_none")]
107 pub trace_id: Option<String>,
108 #[serde(skip_serializing_if = "Option::is_none")]
109 pub task_id: Option<String>,
110 #[serde(skip_serializing_if = "Option::is_none")]
111 pub request_id: Option<String>,
112 #[serde(skip_serializing_if = "Option::is_none")]
113 pub queue_topic: Option<String>,
114 #[serde(skip_serializing_if = "Option::is_none")]
115 pub policy_action: Option<String>,
116 #[serde(skip_serializing_if = "Option::is_none")]
117 pub policy_reason: Option<String>,
118 #[serde(skip_serializing_if = "Option::is_none")]
119 pub retry_count: Option<u32>,
120 #[serde(skip_serializing_if = "Option::is_none")]
121 pub traceback: Option<String>,
122}
123
124impl LogRecord {
125 fn new(level: Level, module: impl Into<String>, message: impl Into<String>) -> Self {
126 let time = Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true);
127 let level_name = level.to_string();
128 Self {
129 time,
130 level,
131 level_name,
132 module: module.into(),
133 message: message.into(),
134 status: None,
135 event_type: None,
136 phase: None,
137 error_kind: None,
138 trace_id: None,
139 task_id: None,
140 request_id: None,
141 queue_topic: None,
142 policy_action: None,
143 policy_reason: None,
144 retry_count: None,
145 traceback: None,
146 }
147 }
148}
149
150pub trait LogSink: Send + Sync {
151 fn name(&self) -> &'static str;
152 fn enabled(&self) -> bool {
153 true
154 }
155 fn min_level(&self) -> Level;
156 fn emit(&self, record: &LogRecord) -> Result<(), LogError>;
157 fn flush(&self) -> Result<(), LogError> {
158 Ok(())
159 }
160}
161
162struct LogDispatcher {
163 sinks: Vec<Arc<dyn LogSink>>,
164}
165
166impl LogDispatcher {
167 fn new(sinks: Vec<Arc<dyn LogSink>>) -> Self {
168 Self { sinks }
169 }
170
171 fn emit(&self, record: LogRecord) {
172 if self.sinks.is_empty() {
173 return;
174 }
175
176 for sink in &self.sinks {
177 if !sink.enabled() || record.level > sink.min_level() {
178 continue;
179 }
180 if sink.emit(&record).is_err() {
181 metrics::counter!("mocra_log_sink_errors_total", "sink" => sink.name())
182 .increment(1);
183 }
184 }
185 }
186}
187
188#[cfg(feature = "dashboard")]
191static LOG_RING: once_cell::sync::Lazy<std::sync::Mutex<std::collections::VecDeque<LogRecord>>> =
192 once_cell::sync::Lazy::new(|| std::sync::Mutex::new(std::collections::VecDeque::new()));
193
194#[cfg(feature = "dashboard")]
195const LOG_RING_CAP: usize = 1000;
196
197#[cfg(feature = "dashboard")]
198fn push_log_ring(record: &LogRecord) {
199 if let Ok(mut ring) = LOG_RING.lock() {
200 if ring.len() >= LOG_RING_CAP {
201 ring.pop_front();
202 }
203 ring.push_back(record.clone());
204 }
205}
206
207#[cfg(feature = "dashboard")]
210pub fn recent_logs(limit: usize) -> Vec<LogRecord> {
211 LOG_RING
212 .lock()
213 .map(|ring| ring.iter().rev().take(limit).cloned().collect())
214 .unwrap_or_default()
215}
216
217struct LogSinkLayer {
218 dispatcher: Arc<LogDispatcher>,
219}
220
221impl LogSinkLayer {
222 fn new(dispatcher: Arc<LogDispatcher>) -> Self {
223 Self { dispatcher }
224 }
225}
226
227impl<S> Layer<S> for LogSinkLayer
228where
229 S: Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>,
230{
231 fn on_event(&self, event: &Event<'_>, _ctx: Context<'_, S>) {
232 let metadata = event.metadata();
233 let mut visitor = LogVisitor::new();
234 event.record(&mut visitor);
235
236 let message = if visitor.message.is_empty() {
237 metadata.name().to_string()
238 } else {
239 visitor.message
240 };
241
242 let mut record = LogRecord::new(*metadata.level(), metadata.target(), message);
243 record.status = visitor.status;
244 record.event_type = visitor.event_type;
245 record.phase = visitor.phase;
246 record.error_kind = visitor.error_kind;
247 record.trace_id = visitor.trace_id;
248 record.task_id = visitor.task_id;
249 record.request_id = visitor.request_id;
250 record.queue_topic = visitor.queue_topic;
251 record.policy_action = visitor.policy_action;
252 record.policy_reason = visitor.policy_reason;
253 record.retry_count = visitor.retry_count;
254 record.traceback = visitor.traceback;
255
256 #[cfg(feature = "dashboard")]
257 push_log_ring(&record);
258 self.dispatcher.emit(record);
259 }
260}
261
262struct ConsoleSink {
263 min_level: Level,
264 writer: Mutex<tracing_appender::non_blocking::NonBlocking>,
265 _guard: WorkerGuard,
266}
267
268impl ConsoleSink {
269 fn new(min_level: Level) -> Self {
270 let (writer, guard) = tracing_appender::non_blocking(std::io::stderr());
271 Self {
272 min_level,
273 writer: Mutex::new(writer),
274 _guard: guard,
275 }
276 }
277}
278
279impl LogSink for ConsoleSink {
280 fn name(&self) -> &'static str {
281 "console"
282 }
283
284 fn min_level(&self) -> Level {
285 self.min_level
286 }
287
288 fn emit(&self, record: &LogRecord) -> Result<(), LogError> {
289 let line = format_log_record_text(record);
290 if let Ok(mut writer) = self.writer.lock() {
291 use std::io::Write;
292 writeln!(writer, "{}", line)?;
293 }
294 metrics::counter!("mocra_log_events_total", "sink" => self.name(), "level" => record.level.as_str()).increment(1);
295 Ok(())
296 }
297}
298
299struct FileSink {
300 min_level: Level,
301 writer: Mutex<tracing_appender::non_blocking::NonBlocking>,
302 _guard: WorkerGuard,
303}
304
305impl FileSink {
306 fn new(path: &Path, min_level: Level, rotation: Rotation) -> Result<Self, LogError> {
307 if let Some(parent) = path.parent() {
308 std::fs::create_dir_all(parent)?;
309 }
310 let file_prefix = path.file_name().and_then(|n| n.to_str()).unwrap_or("app");
311 let file_appender = tracing_appender::rolling::Builder::new()
312 .rotation(rotation)
313 .filename_prefix(file_prefix)
314 .filename_suffix("log")
315 .build(path.parent().unwrap_or_else(|| Path::new(".")))?;
316 let (writer, guard) = tracing_appender::non_blocking(file_appender);
317 Ok(Self {
318 min_level,
319 writer: Mutex::new(writer),
320 _guard: guard,
321 })
322 }
323}
324
325impl LogSink for FileSink {
326 fn name(&self) -> &'static str {
327 "file"
328 }
329
330 fn min_level(&self) -> Level {
331 self.min_level
332 }
333
334 fn emit(&self, record: &LogRecord) -> Result<(), LogError> {
335 let line = format_log_record_text(record);
336 if let Ok(mut writer) = self.writer.lock() {
337 use std::io::Write;
338 writeln!(writer, "{}", line)?;
339 }
340 metrics::counter!("mocra_log_events_total", "sink" => self.name(), "level" => record.level.as_str()).increment(1);
341 Ok(())
342 }
343}
344
345struct DynamicMqSink {
346 min_level: Level,
347}
348
349impl DynamicMqSink {
350 fn new(min_level: Level) -> Self {
351 Self { min_level }
352 }
353}
354
355impl LogSink for DynamicMqSink {
356 fn name(&self) -> &'static str {
357 "mq"
358 }
359
360 fn min_level(&self) -> Level {
361 self.min_level
362 }
363
364 fn emit(&self, record: &LogRecord) -> Result<(), LogError> {
365 let Some(dynamic) = DYNAMIC_SENDER.read().ok().and_then(|g| {
366 g.as_ref()
367 .map(|d| (d.sender.clone(), d.queue_level, d.capacity))
368 }) else {
369 metrics::counter!("mocra_log_dropped_total", "sink" => self.name(), "reason" => "sender_unset").increment(1);
370 return Ok(());
371 };
372
373 if record.level > dynamic.1 {
374 return Ok(());
375 }
376
377 let status = record
378 .status
379 .clone()
380 .or_else(|| record.phase.clone())
381 .unwrap_or_else(|| "info".to_string());
382 let log_model = LogModel {
383 task_id: record
384 .task_id
385 .clone()
386 .unwrap_or_else(|| "unknown".to_string()),
387 request_id: record.request_id.as_ref().and_then(|s| s.parse().ok()),
388 status,
389 level: record.level_name.clone(),
390 message: record.message.clone(),
391 timestamp: record.time.clone(),
392 traceback: record.traceback.clone(),
393 };
394
395 if dynamic.0.try_send(log_model).is_err() {
396 metrics::counter!("mocra_log_dropped_total", "sink" => self.name(), "reason" => "channel_full").increment(1);
397 } else {
398 metrics::counter!("mocra_log_events_total", "sink" => self.name(), "level" => record.level.as_str()).increment(1);
399 metrics::gauge!("mocra_log_batch_size", "sink" => self.name()).set(1.0);
400 }
401
402 if let Some(capacity) = dynamic.2 {
403 let remaining = dynamic.0.capacity();
404 let lag = capacity.saturating_sub(remaining) as f64;
405 metrics::gauge!("mocra_log_queue_lag", "sink" => self.name()).set(lag);
406 }
407 Ok(())
408 }
409}
410
411struct PrometheusSink {
412 min_level: Level,
413}
414
415impl PrometheusSink {
416 fn new(min_level: Level) -> Self {
417 Self { min_level }
418 }
419}
420
421impl LogSink for PrometheusSink {
422 fn name(&self) -> &'static str {
423 "prometheus"
424 }
425
426 fn min_level(&self) -> Level {
427 self.min_level
428 }
429
430 fn emit(&self, record: &LogRecord) -> Result<(), LogError> {
431 metrics::counter!("mocra_log_events_total", "sink" => self.name(), "level" => record.level.as_str()).increment(1);
432 Ok(())
433 }
434}
435
436#[derive(Debug, Clone)]
437pub struct LogSender {
438 pub sender: Sender<LogModel>,
439 pub level: String,
440 pub capacity: Option<usize>,
441}
442
443impl LogSender {
444 pub fn new(sender: Sender<LogModel>, level: impl AsRef<str>) -> Self {
445 Self {
446 sender,
447 level: level.as_ref().into(),
448 capacity: None,
449 }
450 }
451
452 pub fn with_capacity(
453 sender: Sender<LogModel>,
454 level: impl AsRef<str>,
455 capacity: usize,
456 ) -> Self {
457 Self {
458 sender,
459 level: level.as_ref().into(),
460 capacity: Some(capacity),
461 }
462 }
463
464 pub fn with_warn_level(sender: Sender<LogModel>) -> Self {
465 Self::new(sender, "warn")
466 }
467}
468
469struct DynamicSender {
470 sender: Sender<LogModel>,
471 queue_level: Level,
472 capacity: Option<usize>,
473}
474
475static DYNAMIC_SENDER: Lazy<RwLock<Option<DynamicSender>>> = Lazy::new(|| RwLock::new(None));
476
477pub fn set_log_sender(log_sender: LogSender) -> Result<(), Box<dyn std::error::Error>> {
478 let queue_level = log_sender
479 .level
480 .parse::<Level>()
481 .map_err(|_| format!("Invalid queue log level: {}", log_sender.level))?;
482 let mut guard = DYNAMIC_SENDER
483 .write()
484 .expect("DYNAMIC_SENDER write lock poisoned");
485 *guard = Some(DynamicSender {
486 sender: log_sender.sender,
487 queue_level,
488 capacity: log_sender.capacity,
489 });
490 Ok(())
491}
492
493#[allow(dead_code)]
494pub fn clear_log_sender() {
495 if let Ok(mut guard) = DYNAMIC_SENDER.write() {
496 *guard = None;
497 }
498}
499
500#[derive(Debug, Clone)]
501pub enum LogOutputConfig {
502 Console,
503 File {
504 path: PathBuf,
505 rotation: Option<String>,
506 },
507 Mq,
508}
509
510#[derive(Debug, Clone)]
511pub struct PrometheusConfig {
512 pub enabled: bool,
513}
514
515#[derive(Debug, Clone)]
516pub struct LoggerConfig {
517 pub enabled: bool,
518 pub level: String,
519 pub format: String,
520 pub include: Vec<String>,
521 pub buffer: usize,
522 pub flush_interval_ms: u64,
523 pub outputs: Vec<LogOutputConfig>,
524 pub prometheus: Option<PrometheusConfig>,
525}
526
527impl LoggerConfig {
528 pub async fn init(self) -> Result<(), Box<dyn std::error::Error>> {
529 init_logger(self).await
530 }
531
532 pub fn new() -> Self {
533 Self::default()
534 }
535
536 pub fn with_level(mut self, level: impl AsRef<str>) -> Self {
537 self.level = level.as_ref().into();
538 self
539 }
540
541 pub fn with_output(mut self, output: LogOutputConfig) -> Self {
542 self.outputs.push(output);
543 self
544 }
545
546 pub fn for_app(namespace: &str) -> Self {
547 Self {
548 outputs: vec![
549 LogOutputConfig::Console {},
550 LogOutputConfig::File {
551 path: PathBuf::from("logs").join(format!("mocra.{namespace}.log")),
552 rotation: Some("daily".to_string()),
553 },
554 ],
555 ..Self::default()
556 }
557 }
558}
559
560impl Default for LoggerConfig {
561 fn default() -> Self {
562 Self {
563 enabled: true,
564 level: DEFAULT_APP_LOG_LEVEL.to_string(),
565 format: "text".to_string(),
566 include: vec![],
567 buffer: 10000,
568 flush_interval_ms: 500,
569 outputs: vec![LogOutputConfig::Console {}],
570 prometheus: None,
571 }
572 }
573}
574
575const DEFAULT_APP_LOG_LEVEL: &str = "info,engine=debug;sqlx=warn,sea_orm=warn";
576
577static LOGGER_INITIALIZED: AtomicBool = AtomicBool::new(false);
578
579pub fn is_logging_disabled() -> bool {
580 let value = env::var("DISABLE_LOGS")
581 .or_else(|_| env::var("MOCRA_DISABLE_LOGS"))
582 .unwrap_or_default();
583 matches!(
584 value.trim().to_lowercase().as_str(),
585 "1" | "true" | "yes" | "y" | "on"
586 )
587}
588
589pub async fn init_app_logger(namespace: &str) -> Result<bool, Box<dyn std::error::Error>> {
590 if is_logging_disabled() {
591 return Ok(false);
592 }
593
594 let config = LoggerConfig::for_app(namespace);
595 init_logger(config).await?;
596 Ok(true)
597}
598
599pub async fn init_logger(config: LoggerConfig) -> Result<(), Box<dyn std::error::Error>> {
600 if is_logging_disabled() {
601 let _ = LOGGER_INITIALIZED.swap(true, Ordering::SeqCst);
602 return Ok(());
603 }
604 if LOGGER_INITIALIZED.swap(true, Ordering::SeqCst) {
605 tracing::warn!("Logger already initialized, skipping re-initialization");
606 return Ok(());
607 }
608
609 let _ = LogTracer::builder()
610 .with_max_level(log::LevelFilter::Trace)
611 .init();
612
613 let configured_filter = normalize_filter_string(&config.level);
614 let filter = if configured_filter != DEFAULT_APP_LOG_LEVEL {
615 EnvFilter::try_new(&configured_filter).unwrap_or_else(|_| EnvFilter::new("info"))
616 } else {
617 EnvFilter::try_from_default_env()
618 .or_else(|_| EnvFilter::try_new(&configured_filter))
619 .unwrap_or_else(|_| EnvFilter::new("info"))
620 };
621
622 let sinks = build_sinks(&config)?;
623 let dispatcher = Arc::new(LogDispatcher::new(sinks));
624 let layer = LogSinkLayer::new(dispatcher);
625
626 let _ = tracing_subscriber::registry()
627 .with(layer)
628 .with(filter)
629 .try_init();
630
631 Ok(())
632}
633
634pub async fn init_simple_logger() -> Result<(), Box<dyn std::error::Error>> {
635 let config = LoggerConfig::default();
636 init_logger(config).await
637}
638
639fn build_sinks(config: &LoggerConfig) -> Result<Vec<Arc<dyn LogSink>>, LogError> {
640 if !config.enabled {
641 return Ok(Vec::new());
642 }
643
644 let mut sinks: Vec<Arc<dyn LogSink>> = Vec::new();
645 let base_level = base_level_from_filter(&config.level).unwrap_or(Level::INFO);
646
647 for output in &config.outputs {
648 match output {
649 LogOutputConfig::Console => {
650 sinks.push(Arc::new(ConsoleSink::new(base_level)));
651 }
652 LogOutputConfig::File { path, rotation } => {
653 let rotation = match rotation.as_deref() {
654 Some("daily") | None => Rotation::DAILY,
655 Some("hourly") => Rotation::HOURLY,
656 Some("never") => Rotation::NEVER,
657 Some("minutely") => Rotation::MINUTELY,
658 _ => Rotation::DAILY,
659 };
660 sinks.push(Arc::new(FileSink::new(
661 path.as_path(),
662 base_level,
663 rotation,
664 )?));
665 }
666 LogOutputConfig::Mq => {
667 sinks.push(Arc::new(DynamicMqSink::new(base_level)));
668 }
669 }
670 }
671
672 if let Some(prometheus) = &config.prometheus
673 && prometheus.enabled
674 {
675 sinks.push(Arc::new(PrometheusSink::new(base_level)));
676 }
677
678 Ok(sinks)
679}
680
681fn normalize_filter_string(filter: &str) -> String {
682 let trimmed = filter.trim();
683 if trimmed.contains('=') || trimmed.contains(',') || trimmed.contains(';') {
684 return trimmed.to_string();
685 }
686 let lower = trimmed.to_lowercase();
687 let normalized = match lower.as_str() {
688 "all" => "trace",
689 "fatal" => "error",
690 "warning" => "warn",
691 other => other,
692 };
693 build_allowlist_filter(normalized)
694}
695
696fn build_allowlist_filter(level: &str) -> String {
697 format!(
698 "off,cacheable={level},common={level},downloader={level},engine={level},errors={level},js_v8={level},mocra={level},proxy={level},queue={level},sync={level},utils={level},tests={level},python_mocra={level},sqlx=warn,sea_orm=warn"
699 )
700}
701
702fn base_level_from_filter(level: &str) -> Option<Level> {
703 let candidate = level
704 .split([',', ';'])
705 .next()
706 .map(|value| value.trim())
707 .filter(|value| !value.is_empty())?;
708 candidate.parse::<Level>().ok()
709}
710
711fn format_log_record_text(record: &LogRecord) -> String {
712 let mut line = format!(
713 "{} [{}] {} - {}",
714 record.time, record.level_name, record.module, record.message
715 );
716
717 if let Some(value) = &record.event_type {
718 line.push_str(&format!(" event_type={value}"));
719 }
720 if let Some(value) = &record.status {
721 line.push_str(&format!(" status={value}"));
722 }
723 if let Some(value) = &record.phase {
724 line.push_str(&format!(" phase={value}"));
725 }
726 if let Some(value) = &record.error_kind {
727 line.push_str(&format!(" error_kind={value}"));
728 }
729 if let Some(value) = &record.trace_id {
730 line.push_str(&format!(" trace_id={value}"));
731 }
732 if let Some(value) = &record.task_id {
733 line.push_str(&format!(" task_id={value}"));
734 }
735 if let Some(value) = &record.request_id {
736 line.push_str(&format!(" request_id={value}"));
737 }
738 if let Some(value) = &record.queue_topic {
739 line.push_str(&format!(" queue_topic={value}"));
740 }
741 if let Some(value) = &record.policy_action {
742 line.push_str(&format!(" policy.action={value}"));
743 }
744 if let Some(value) = &record.policy_reason {
745 line.push_str(&format!(" policy.reason={value}"));
746 }
747 if let Some(value) = &record.retry_count {
748 line.push_str(&format!(" retry.count={value}"));
749 }
750 if let Some(value) = &record.traceback {
751 line.push_str(&format!(" traceback={value}"));
752 }
753
754 line
755}
756
757struct LogVisitor {
758 message: String,
759 status: Option<String>,
760 event_type: Option<String>,
761 phase: Option<String>,
762 error_kind: Option<String>,
763 trace_id: Option<String>,
764 task_id: Option<String>,
765 request_id: Option<String>,
766 queue_topic: Option<String>,
767 policy_action: Option<String>,
768 policy_reason: Option<String>,
769 retry_count: Option<u32>,
770 traceback: Option<String>,
771}
772
773impl LogVisitor {
774 fn new() -> Self {
775 Self {
776 message: String::with_capacity(64),
777 status: None,
778 event_type: None,
779 phase: None,
780 error_kind: None,
781 trace_id: None,
782 task_id: None,
783 request_id: None,
784 queue_topic: None,
785 policy_action: None,
786 policy_reason: None,
787 retry_count: None,
788 traceback: None,
789 }
790 }
791}
792
793impl Visit for LogVisitor {
794 fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
795 if field.name() == "message" {
796 use std::fmt::Write;
797 let _ = write!(self.message, "{:?}", value);
798 }
799 }
800
801 fn record_str(&mut self, field: &Field, value: &str) {
802 match field.name() {
803 "message" => self.message.push_str(value),
804 "status" => self.status = Some(value.to_string()),
805 "event_type" => self.event_type = Some(value.to_string()),
806 "phase" => self.phase = Some(value.to_string()),
807 "error_kind" => self.error_kind = Some(value.to_string()),
808 "trace_id" => self.trace_id = Some(value.to_string()),
809 "task_id" => self.task_id = Some(value.to_string()),
810 "request_id" => self.request_id = Some(value.to_string()),
811 "queue_topic" => self.queue_topic = Some(value.to_string()),
812 "policy_action" => self.policy_action = Some(value.to_string()),
813 "policy_reason" => self.policy_reason = Some(value.to_string()),
814 "traceback" => self.traceback = Some(value.to_string()),
815 _ => {}
816 }
817 }
818
819 fn record_u64(&mut self, field: &Field, value: u64) {
820 if field.name() == "retry_count" {
821 self.retry_count = Some(value as u32);
822 }
823 }
824}
825
826#[cfg(test)]
827mod tests {
828 use super::*;
829 use tokio::sync::mpsc;
830 use tracing::{debug, error, info, warn};
831
832 #[test]
833 fn test_logger_config_builder() {
834 let config = LoggerConfig::new()
835 .with_level("debug")
836 .with_output(LogOutputConfig::Console {});
837
838 assert_eq!(config.level, "debug");
839 assert!(!config.outputs.is_empty());
840 }
841
842 #[tokio::test]
843 async fn test_simple_logger_init() {
844 let config = LoggerConfig::new().with_level("info");
845 let _ = init_logger(config).await;
846 }
847
848 #[tokio::test]
849 async fn test_log_levels() {
850 let config = LoggerConfig::new().with_level("debug");
851 let _ = init_logger(config).await;
852
853 debug!("Debug message");
854 info!("Info message");
855 warn!("Warning message");
856 error!("Error message");
857 }
858
859 #[tokio::test]
860 async fn test_queue_logger() {
861 let (sender, mut receiver) = mpsc::channel(100);
862 let log_sender = LogSender::new(sender, "warn");
863 let _ = set_log_sender(log_sender);
864
865 let config = LoggerConfig::new()
866 .with_level("info")
867 .with_output(LogOutputConfig::Mq {});
868
869 let _ = init_logger(config).await;
870
871 info!(
872 task_id = "test-task",
873 status = "success",
874 "Info message should not go to queue"
875 );
876 warn!(
877 task_id = "test-task",
878 status = "warning",
879 "Warning message for queue"
880 );
881
882 tokio::time::sleep(tokio::time::Duration::from_millis(100)).await;
883
884 if let Ok(log_model) = receiver.try_recv() {
885 assert_eq!(log_model.task_id, "test-task");
886 assert_eq!(log_model.status, "warning");
887 assert_eq!(log_model.level, "WARN");
888 assert!(log_model.message.contains("Warning message for queue"));
889 }
890
891 assert!(receiver.try_recv().is_err());
892 }
893}