use crate::LogRecord;
use crate::Metrics;
use chrono::Utc;
use crossbeam_channel::Sender;
use log::{Level, LevelFilter, Metadata, Record};
use std::sync::Arc;
pub struct LogAdapter {
console_sender: Sender<Arc<LogRecord>>,
async_sender: Sender<Arc<LogRecord>>,
metrics: Arc<Metrics>,
}
impl LogAdapter {
pub fn new(
console_sender: Sender<Arc<LogRecord>>,
async_sender: Sender<Arc<LogRecord>>,
metrics: Arc<Metrics>,
) -> Self {
Self {
console_sender,
async_sender,
metrics,
}
}
fn level_to_string(level: Level) -> &'static str {
match level {
Level::Trace => "TRACE",
Level::Debug => "DEBUG",
Level::Info => "INFO",
Level::Warn => "WARN",
Level::Error => "ERROR",
}
}
fn record_to_log_record(&self, record: &Record) -> LogRecord {
LogRecord {
timestamp: Utc::now(),
level: Self::level_to_string(record.level()).to_string(),
target: record.target().to_string(),
message: record.args().to_string(),
file: record.file().map(|s| s.to_string()),
line: record.line(),
thread_id: std::thread::current()
.name()
.unwrap_or("unknown")
.to_string(),
fields: Default::default(),
}
}
}
impl log::Log for LogAdapter {
fn enabled(&self, metadata: &Metadata) -> bool {
metadata.level() <= log::max_level()
}
fn log(&self, record: &Record) {
if !self.enabled(record.metadata()) {
return;
}
let log_record = Arc::new(self.record_to_log_record(record));
match self.console_sender.try_send(Arc::clone(&log_record)) {
Ok(_) => {}
Err(crossbeam_channel::TrySendError::Full(_)) => {
self.metrics.inc_channel_blocked();
self.metrics.inc_logs_dropped();
}
Err(crossbeam_channel::TrySendError::Disconnected(_)) => {
self.metrics.inc_logs_dropped();
}
}
match self.async_sender.try_send(log_record) {
Ok(_) => {}
Err(crossbeam_channel::TrySendError::Full(_)) => {
self.metrics.inc_channel_blocked();
self.metrics.inc_logs_dropped();
}
Err(crossbeam_channel::TrySendError::Disconnected(_)) => {
self.metrics.inc_logs_dropped();
}
}
}
fn flush(&self) {
}
}
pub struct LogLogger {
adapter: LogAdapter,
max_level: LevelFilter,
}
impl LogLogger {
pub fn new(adapter: LogAdapter, max_level: LevelFilter) -> Self {
Self { adapter, max_level }
}
pub fn install(self) -> Result<(), log::SetLoggerError> {
let max_level = self.max_level;
log::set_boxed_logger(Box::new(self))?;
log::set_max_level(max_level);
Ok(())
}
}
impl log::Log for LogLogger {
fn enabled(&self, metadata: &Metadata) -> bool {
self.adapter.enabled(metadata)
}
fn log(&self, record: &Record) {
self.adapter.log(record)
}
fn flush(&self) {
self.adapter.flush()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crossbeam_channel::bounded;
use log::Log;
#[test]
fn test_level_to_string() {
assert_eq!(LogAdapter::level_to_string(Level::Error), "ERROR");
assert_eq!(LogAdapter::level_to_string(Level::Warn), "WARN");
assert_eq!(LogAdapter::level_to_string(Level::Info), "INFO");
assert_eq!(LogAdapter::level_to_string(Level::Debug), "DEBUG");
assert_eq!(LogAdapter::level_to_string(Level::Trace), "TRACE");
}
#[test]
fn test_record_to_log_record() {
let (console_tx, _) = bounded(100);
let (async_tx, _) = bounded(100);
let metrics = Arc::new(Metrics::new());
let adapter = LogAdapter::new(console_tx, async_tx, metrics);
let metadata = log::Metadata::builder()
.target("test::module")
.level(Level::Info)
.build();
let record = log::Record::builder()
.metadata(metadata)
.args(format_args!("Test message"))
.file(Some("test.rs"))
.line(Some(42))
.build();
let log_record = adapter.record_to_log_record(&record);
assert_eq!(log_record.level, "INFO");
assert_eq!(log_record.target, "test::module");
assert_eq!(log_record.message, "Test message");
assert_eq!(log_record.file, Some("test.rs".to_string()));
assert_eq!(log_record.line, Some(42));
}
#[test]
fn test_log_adapter_log_sends_to_channels() {
let (console_tx, console_rx) = bounded(10);
let (async_tx, async_rx) = bounded(10);
let metrics = Arc::new(Metrics::new());
let adapter = LogAdapter::new(console_tx, async_tx, metrics);
log::set_max_level(log::LevelFilter::Info);
let metadata = log::Metadata::builder()
.target("test::adapter")
.level(Level::Info)
.build();
let record = log::Record::builder()
.metadata(metadata)
.args(format_args!("Adapter send"))
.file(Some("test.rs"))
.line(Some(7))
.build();
adapter.log(&record);
let console_received = console_rx.recv().unwrap();
assert_eq!(console_received.level, "INFO");
assert_eq!(console_received.target, "test::adapter");
assert_eq!(console_received.message, "Adapter send");
let async_received = async_rx.recv().unwrap();
assert_eq!(async_received.level, "INFO");
assert_eq!(async_received.target, "test::adapter");
assert_eq!(async_received.message, "Adapter send");
}
#[test]
fn test_log_adapter_handles_full_channel() {
let (console_tx, console_rx) = bounded(1);
let (async_tx, async_rx) = bounded(1);
let metrics = Arc::new(Metrics::new());
let adapter = LogAdapter::new(console_tx, async_tx, metrics.clone());
log::set_max_level(log::LevelFilter::Info);
for i in 0..5 {
let metadata = log::Metadata::builder()
.target("test::adapter")
.level(Level::Info)
.build();
let msg = format!("Test message {}", i);
let args = format_args!("{}", msg);
let record = log::Record::builder().metadata(metadata).args(args).build();
adapter.log(&record);
}
while console_rx.try_recv().is_ok() {}
while async_rx.try_recv().is_ok() {}
assert_eq!(metrics.logs_dropped(), 8);
}
#[test]
fn test_log_adapter_disconnected_channel() {
let (console_tx, _cr) = bounded(10);
let (async_tx, _ar) = bounded(1);
drop(_ar); let metrics = Arc::new(Metrics::new());
let adapter = LogAdapter::new(console_tx, async_tx, metrics.clone());
log::set_max_level(log::LevelFilter::Info);
let metadata = log::Metadata::builder()
.target("test::adapter")
.level(Level::Info)
.build();
let record = log::Record::builder()
.metadata(metadata)
.args(format_args!("after disconnect"))
.build();
adapter.log(&record);
assert_eq!(metrics.logs_dropped(), 1);
}
#[test]
fn test_flush_is_noop() {
let (console_tx, _) = bounded(10);
let (async_tx, _) = bounded(10);
let metrics = Arc::new(Metrics::new());
let adapter = LogAdapter::new(console_tx, async_tx, metrics);
adapter.flush();
}
#[test]
fn test_log_adapter_all_levels_mapped() {
let (console_tx, _) = bounded(10);
let (async_tx, _) = bounded(10);
let metrics = Arc::new(Metrics::new());
let adapter = LogAdapter::new(console_tx, async_tx, metrics.clone());
log::set_max_level(log::LevelFilter::Trace);
for (level, expected_str) in [
(Level::Error, "ERROR"),
(Level::Warn, "WARN"),
(Level::Info, "INFO"),
(Level::Debug, "DEBUG"),
(Level::Trace, "TRACE"),
] {
let metadata = log::Metadata::builder()
.target("test::levels")
.level(level)
.build();
let args = format_args!("msg for {expected_str}");
let record = log::Record::builder().metadata(metadata).args(args).build();
let log_record = adapter.record_to_log_record(&record);
assert_eq!(
log_record.level, expected_str,
"level {level:?} should map to {expected_str}"
);
}
}
#[test]
fn test_log_adapter_enabled_respects_max_level() {
let (console_tx, _) = bounded(10);
let (async_tx, _) = bounded(10);
let metrics = Arc::new(Metrics::new());
let adapter = LogAdapter::new(console_tx, async_tx, metrics);
log::set_max_level(log::LevelFilter::Info);
let info_meta = log::Metadata::builder()
.target("test")
.level(Level::Info)
.build();
let debug_meta = log::Metadata::builder()
.target("test")
.level(Level::Debug)
.build();
let trace_meta = log::Metadata::builder()
.target("test")
.level(Level::Trace)
.build();
assert!(adapter.enabled(&info_meta));
assert!(!adapter.enabled(&debug_meta));
assert!(!adapter.enabled(&trace_meta));
}
#[test]
fn test_log_adapter_console_disconnected_channel() {
let (console_tx, _cr) = bounded(10);
let (async_tx, _ar) = bounded(10);
drop(_cr); drop(_ar); let metrics = Arc::new(Metrics::new());
let adapter = LogAdapter::new(console_tx, async_tx, metrics.clone());
log::set_max_level(log::LevelFilter::Info);
let metadata = log::Metadata::builder()
.target("test::adapter")
.level(Level::Info)
.build();
let record = log::Record::builder()
.metadata(metadata)
.args(format_args!("console disconnected"))
.build();
adapter.log(&record);
assert_eq!(metrics.logs_dropped(), 2);
}
#[test]
fn test_log_logger_new() {
let (console_tx, _) = bounded(10);
let (async_tx, _) = bounded(10);
let metrics = Arc::new(Metrics::new());
let adapter = LogAdapter::new(console_tx, async_tx, metrics);
let logger = LogLogger::new(adapter, LevelFilter::Info);
let _ = logger;
}
#[test]
fn test_log_logger_enabled() {
let (console_tx, _) = bounded(10);
let (async_tx, _) = bounded(10);
let metrics = Arc::new(Metrics::new());
let adapter = LogAdapter::new(console_tx, async_tx, metrics);
let logger = LogLogger::new(adapter, LevelFilter::Info);
log::set_max_level(log::LevelFilter::Info);
let info_meta = log::Metadata::builder()
.target("test")
.level(Level::Info)
.build();
let debug_meta = log::Metadata::builder()
.target("test")
.level(Level::Debug)
.build();
assert!(logger.enabled(&info_meta));
assert!(!logger.enabled(&debug_meta));
}
#[test]
fn test_log_logger_log() {
let (console_tx, console_rx) = bounded(10);
let (async_tx, async_rx) = bounded(10);
let metrics = Arc::new(Metrics::new());
let adapter = LogAdapter::new(console_tx, async_tx, metrics);
let logger = LogLogger::new(adapter, LevelFilter::Info);
log::set_max_level(log::LevelFilter::Info);
let metadata = log::Metadata::builder()
.target("test::logger")
.level(Level::Info)
.build();
let record = log::Record::builder()
.metadata(metadata)
.args(format_args!("via LogLogger"))
.build();
logger.log(&record);
let console_received = console_rx.recv().unwrap();
assert_eq!(console_received.message, "via LogLogger");
let async_received = async_rx.recv().unwrap();
assert_eq!(async_received.message, "via LogLogger");
}
#[test]
fn test_log_logger_flush() {
let (console_tx, _) = bounded(10);
let (async_tx, _) = bounded(10);
let metrics = Arc::new(Metrics::new());
let adapter = LogAdapter::new(console_tx, async_tx, metrics);
let logger = LogLogger::new(adapter, LevelFilter::Info);
logger.flush();
}
#[test]
#[serial_test::serial]
fn test_log_logger_install_second_call_err_after_first() {
let (console_tx1, _cr1) = bounded(10);
let (async_tx1, _ar1) = bounded(10);
let metrics1 = Arc::new(Metrics::new());
let adapter1 = LogAdapter::new(console_tx1, async_tx1, metrics1);
let logger1 = LogLogger::new(adapter1, LevelFilter::Info);
let result1 = logger1.install();
let (console_tx2, _cr2) = bounded(10);
let (async_tx2, _ar2) = bounded(10);
let metrics2 = Arc::new(Metrics::new());
let adapter2 = LogAdapter::new(console_tx2, async_tx2, metrics2);
let logger2 = LogLogger::new(adapter2, LevelFilter::Info);
let result2 = logger2.install();
assert!(
result2.is_err(),
"second install should fail because a global logger is already installed, \
got: {:?} (first install result: {:?})",
result2,
result1
);
}
}