Skip to main content

mocra_core/utils/
logger.rs

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/// Dashboard log ring buffer: captures recent logs for `GET /observability/logs` (`dashboard`
189/// feature only).
190#[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/// The most recent `limit` log records (newest first). Used by the dashboard's
208/// `GET /observability/logs`.
209#[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}