log-full 0.0.1

A simple, asynchronous log library
Documentation
//! 核心日志器模块
//! 
//! 包含 LogPlus 日志器的实现和相关数据结构


use crate::utils::{AsyncLogType, get_msg_from_cache, put_msg_to_cache, level_color};
use crate::quickwit::{QuickwitClient, QuickwitLogEntry};
use std::collections::{HashMap, VecDeque};
use std::fs::File;
use std::io::{LineWriter, Stdout, Write};
use std::sync::{Mutex, RwLock, mpsc::Sender};

/// 自定义过滤函数 trait
pub trait CustomFilter: Send + Sync + 'static {
    fn enabled(&self, record: &log::Record) -> bool;
}

/// 日志器数据
pub struct LogData {
    pub log_size: u32,
    pub console: Option<LineWriter<Stdout>>,
    pub fileout: Option<LineWriter<File>>,
    pub sender: Option<Sender<AsyncLogType>>,
    pub plugin: Option<Box<dyn std::io::Write + Send + Sync + 'static>>,
}

/// 核心日志器
pub struct LogPlus {
    pub level: log::LevelFilter,
    pub log_file: String,
    pub max_size: u32,
    pub level_filter: RwLock<HashMap<String, log::LevelFilter>>,
    pub fmt_cache: Mutex<VecDeque<Vec<u8>>>,
    pub filter: Option<Box<dyn CustomFilter>>,
    pub logger_data: Mutex<LogData>,
    pub show_process_id: bool,
    pub show_thread_info: bool,
    pub show_module_path: bool,
    pub highlight_keywords: Vec<String>,
    pub quickwit_client: Option<QuickwitClient>,
}

impl LogPlus {
    /// 创建新的日志器实例
    pub fn new(
        level: log::LevelFilter,
        log_file: String,
        max_size: u32,
        filter: Option<Box<dyn CustomFilter>>,
        logger_data: LogData,
        show_process_id: bool,
        show_thread_info: bool,
        show_module_path: bool,
        highlight_keywords: Vec<String>,
        quickwit_client: Option<QuickwitClient>,
    ) -> Self {
        Self {
            level,
            log_file,
            max_size,
            level_filter: RwLock::new(HashMap::new()),
            fmt_cache: Mutex::new(VecDeque::new()),
            filter,
            logger_data: Mutex::new(logger_data),
            show_process_id,
            show_thread_info,
            show_module_path,
            highlight_keywords,
            quickwit_client,
        }
    }

    /// 输出日志到控制台和文件
    pub fn write(&self, msg: &[u8]) {
        let mut logger_data = match self.logger_data.lock() {
            Ok(v) => v,
            Err(e) => {
                eprint!("log mutex lock failed: {e:?}");
                return;
            }
        };

        // 如果启用了控制台输出,则写入控制台
        if let Some(ref mut console) = logger_data.console {
            console.write_all(msg).expect("write log to console fail");
        }

        // 判断日志长度是否到达最大限制,如果到了,需要备份当前日志文件并重新创建新的日志文件
        if logger_data.log_size > self.max_size {
            let mut log_file_closed = false;

            // 如果启用了日志文件,刷新缓存并关闭日志文件
            if let Some(ref mut fileout) = logger_data.fileout {
                fileout.flush().expect("flush log file fail");
                logger_data.fileout.take();
                log_file_closed = true;
            }

            // 之所以把关闭文件和重新创建文件分开写,是因为rust限制了可变借用(fileout)只允许1次
            if log_file_closed {
                // 删除已有备份,并重命名现有文件为备份文件
                let bak = format!("{}.bak", self.log_file);
                std::fs::remove_file(&bak).unwrap_or_default();
                std::fs::rename(&self.log_file, &bak).expect("backup log file fail");

                match crate::utils::open_log_file_sync(&self.log_file) {
                    Ok((writer, _)) => {
                        logger_data.fileout = Some(writer);
                    },
                    Err(e) => {
                        eprintln!("Failed to reopen log file {}: {}", self.log_file, e);
                        return;
                    }
                }
                logger_data.log_size = 0;
            }
        }

        if let Some(ref mut fileout) = logger_data.fileout {
            let ws = crate::utils::write_text(fileout, msg).unwrap();
            logger_data.log_size += ws as u32;
        }

        if let Some(plugin) = &mut logger_data.plugin {
            plugin.write_all(msg).expect("write log to plugin fail");
        }
    }

    /// 刷新日志的控制台和文件缓存
    pub fn flush_inner(&self) {
        let mut logger_data = match self.logger_data.lock() {
            Ok(v) => v,
            Err(e) => {
                eprint!("log mutex lock failed: {e:?}");
                return;
            }
        };

        if let Some(ref mut console) = logger_data.console {
            if let Err(e) = crate::utils::safe_flush(console) {
                eprintln!("Failed to flush console: {}", e);
            }
        }

        if let Some(ref mut fileout) = logger_data.fileout {
            if let Err(e) = crate::utils::safe_flush(fileout) {
                eprintln!("Failed to flush log file: {}", e);
            }
        }
    }

    /// 应用关键词高亮
    fn apply_keyword_highlighting(&self, message: &str) -> String {
        if self.highlight_keywords.is_empty() {
            return message.to_string();
        }
        
        let mut result = message.to_string();
        for keyword in &self.highlight_keywords {
            if message.contains(keyword) {
                let highlighted = format!("\x1b[41m{}\x1b[0m", keyword);
                result = result.replace(keyword, &highlighted);
            }
        }
        result
    }
}

impl log::Log for LogPlus {
    fn enabled(&self, metadata: &log::Metadata) -> bool {
        if metadata.level() <= self.level {
            if let Ok(level_filters) = self.level_filter.read() {
                let mut target = metadata.target();
                while !target.is_empty() {
                    if let Some(level) = level_filters.get(target) {
                        return metadata.level() <= *level;
                    }

                    target = match target.rfind("::") {
                        Some(rpos) => &target[..rpos],
                        None => ""
                    };
                }
                return true;
            }
        }
        false
    }

    fn log(&self, record: &log::Record) {
        if !self.enabled(record.metadata()) { return; }
        if let Some(filter) = &self.filter {
            if !filter.enabled(record) { return; }
        }

        let now = std::time::SystemTime::now()
            .duration_since(std::time::UNIX_EPOCH)
            .unwrap()
            .as_secs();
        let now = format!("{}", now);

        let mut msg = get_msg_from_cache();

        // 日志条目格式化 - 优化后的格式
        let is_detail = self.level >= log::LevelFilter::Debug;
        let log_level = record.level();
        
        if is_detail {
            // 详细模式:包含颜色、位置信息等
            write!(&mut msg, "\x1b[90m{}\x1b[0m {}{:>5}\x1b[0m",
                now, level_color(log_level), log_level
            ).unwrap();
            
            // 根据配置添加进程ID
            if self.show_process_id {
                write!(&mut msg, " \x1b[90m[{}]\x1b[0m", std::process::id()).unwrap();
            }
            
            // 根据配置添加线程信息
            if self.show_thread_info {
                if let Some(thread_name) = std::thread::current().name() {
                    write!(&mut msg, " \x1b[90m[{}]\x1b[0m", thread_name).unwrap();
                } else {
                    write!(&mut msg, " \x1b[90m[{:?}]\x1b[0m", std::thread::current().id()).unwrap();
                }
            }
            
            // 根据配置添加模块和行号信息
            if self.show_module_path {
                let target = record.target();
                let line = record.line().unwrap_or(0);
                write!(&mut msg, " \x1b[90m{}:{}\x1b[0m", target, line).unwrap();
            }
            
            // 消息内容 - 应用关键词高亮
            let message = format!("{}", record.args());
            let highlighted_message = self.apply_keyword_highlighting(&message);
            write!(&mut msg, "{}\n", highlighted_message).unwrap();
        } else {
            // 简洁模式:只包含基本信息 - 应用关键词高亮
            let message = format!("{}", record.args());
            let highlighted_message = self.apply_keyword_highlighting(&message);
            write!(&mut msg, "{} {:>5}{}\n", now, log_level, highlighted_message).unwrap();
        }

        // 异步写入模式
        match self.logger_data.lock() {
            Ok(logger_data) => {
                if let Some(ref sender) = logger_data.sender {
                    // 采用独立的单线程写入日志的方式,向channel发送要写入的日志消息即可
                    sender.send(AsyncLogType::Message(msg)).unwrap();
                    
                    // 发送到 Quickwit(如果配置了)- 异步模式
                    if let Some(ref _quickwit_client) = self.quickwit_client {
                        let log_entry = QuickwitLogEntry {
                            timestamp: std::time::SystemTime::now()
                                .duration_since(std::time::UNIX_EPOCH)
                                .unwrap()
                                .as_secs(),
                            level: record.level().to_string(),
                            message: format!("{}", record.args()),
                            module: record.module_path().map(|s| s.to_string()),
                            file: record.file().map(|s| s.to_string()),
                            line: record.line(),
                            process_id: Some(std::process::id()),
                            thread_id: std::thread::current().name().map(|s| s.to_string()),
                            custom_fields: std::collections::HashMap::new(),
                        };
                        
                        // 发送 Quickwit 日志到异步处理器
                        if let Err(e) = sender.send(AsyncLogType::QuickwitLog(log_entry)) {
                            eprintln!("✗ 发送 Quickwit 日志到异步处理器失败: {}", e);
                        }
                    }
                    return;
                }
            },
            Err(e) => {
                eprint!("log mutex lock failed: {e:?}");
                return;
            }
        }

        // 同步写入模式
        self.write(&msg);
        put_msg_to_cache(msg.clone());
        
        // 发送到 Quickwit(如果配置了)
          if let Some(ref quickwit_client) = self.quickwit_client {
              let log_entry = QuickwitLogEntry {
                  timestamp: std::time::SystemTime::now()
                      .duration_since(std::time::UNIX_EPOCH)
                      .unwrap()
                      .as_secs(),
                  level: record.level().to_string(),
                  message: format!("{}", record.args()),
                  module: record.module_path().map(|s| s.to_string()),
                  file: record.file().map(|s| s.to_string()),
                  line: record.line(),
                  process_id: Some(std::process::id()),
                  thread_id: std::thread::current().name().map(|s| s.to_string()),
                  custom_fields: std::collections::HashMap::new(),
              };
              
              // 同步发送到 Quickwit,避免生命周期问题
              match quickwit_client.send_log(&log_entry) {
                  Ok(()) => {
                      // 可选:添加调试信息
                      if std::env::var("QUICKWIT_DEBUG").is_ok() {
                          eprintln!("✓ 日志已发送到 Quickwit: {}", log_entry.message);
                      }
                  },
                  Err(e) => {
                      eprintln!("✗ 发送日志到 Quickwit 失败: {}", e);
                      eprintln!("  日志内容: {}", log_entry.message);
                      eprintln!("  错误详情: {:?}", e);
                  }
              }
          }
    }

    fn flush(&self) {
        if let Ok(logger_data) = self.logger_data.lock() {
            if let Some(ref sender) = logger_data.sender {
                if let Err(e) = sender.send(AsyncLogType::Flush) {
                    eprint!("failed in log::flush: {e:?}");
                }
            } else {
                drop(logger_data);
                self.flush_inner();
            }
        }
    }
}

/// 为函数类型实现 CustomFilter trait
impl<F: Fn(&log::Record) -> bool + Send + Sync + 'static> CustomFilter for F {
    fn enabled(&self, record: &log::Record) -> bool {
        self(record)
    }
}