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};
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;
}
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();
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 {
sender.send(AsyncLogType::Message(msg)).unwrap();
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(),
};
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());
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(),
};
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();
}
}
}
}
impl<F: Fn(&log::Record) -> bool + Send + Sync + 'static> CustomFilter for F {
fn enabled(&self, record: &log::Record) -> bool {
self(record)
}
}