pub(crate) mod mongodb;
pub mod operation;
use std::net::SocketAddr;
use std::rc::Rc;
use std::str::FromStr;
use std::{io, sync::RwLock};
use log::{Level, LevelFilter, Record};
use log4rs::append::console::{ConsoleAppender, Target};
use log4rs::append::rolling_file::policy::compound::roll::fixed_window::FixedWindowRoller;
use log4rs::append::rolling_file::policy::compound::trigger::size::SizeTrigger;
use log4rs::append::rolling_file::policy::compound::trigger::Trigger;
use log4rs::append::rolling_file::policy::compound::CompoundPolicy;
use log4rs::append::rolling_file::RollingFileAppender;
use log4rs::config::{Appender, Logger, Root};
use log4rs::encode::pattern::PatternEncoder;
use log4rs::filter::threshold::ThresholdFilter;
use log4rs::Config;
use tracing_subscriber::layer::SubscriberExt;
use crate::tina::constant::Constants;
use crate::tina::data::AppResult;
use crate::tina::log::log_api::Metadata;
use crate::tina::util::Utility;
use crate::{app_error_from, app_error_from_none_static, app_system_error};
use chrono::Local;
use indexmap::IndexMap;
pub use log as log_api;
pub use log4rs as log_impl;
use log4rs::append::Append;
use std::sync::mpsc::SyncSender;
use std::sync::{Arc, Mutex};
use std::thread::current;
use super::util::string_joiner::StringJoiner;
#[allow(unused)]
static GLOBAL_CLOSEABLE_APPENDERS: Mutex<Vec<Arc<dyn CloseableAppender>>> = Mutex::new(Vec::new());
#[allow(dead_code)]
pub(crate) fn wait_for_closeable_appender_shutdown() -> AppResult<()> {
std::thread::spawn(move || {
let lock = GLOBAL_CLOSEABLE_APPENDERS.lock().map_err(app_error_from_none_static!())?;
for appender in lock.iter() {
appender.shutdown();
}
AppResult::Ok(())
})
.join()
.map_err(|_| app_system_error!("wait_for_closeable_appender_shutdown failed"))?
}
pub trait CloseableAppender: Append + Send + Sync + 'static {
fn shutdown(&self);
}
#[derive(Copy, Clone)]
pub struct MaxFileSize(pub usize);
#[derive(Copy, Clone)]
pub struct MaxFileHistory(pub usize);
pub enum WriteTarget {
Console(String, Target, LevelFilter),
RollingFile(String, String, LevelFilter, MaxFileSize, MaxFileHistory),
#[cfg(feature = "file-store-mongodb")]
Mongodb(String, String, ::mongodb::Database, LevelFilter, std::time::Duration),
}
struct CustomLogConfig(pub String, pub LevelFilter, pub Vec<String>, pub bool);
pub struct LogConfig {
log_home: String,
web_server_name: String,
targets: Vec<WriteTarget>,
root_config: Vec<String>,
root_level: LevelFilter,
custom_config: IndexMap<String, CustomLogConfig>,
monitoring_address: Option<SocketAddr>,
enable_span: bool,
}
impl LogConfig {
pub const DEFAULT_LOG_FILE_SIZE: usize = 300 * 1024 * 1024;
pub fn kb(v: usize) -> usize {
v * 1024
}
pub fn mb(v: usize) -> usize {
v * 1024 * 1024
}
pub fn gb(v: usize) -> usize {
v * 1024 * 1024 * 1024
}
pub fn new(log_home: &str, web_server_name: &str) -> LogConfig {
LogConfig {
log_home: log_home.to_string(),
web_server_name: web_server_name.to_string(),
targets: vec![],
root_config: vec![],
root_level: LevelFilter::Debug,
custom_config: IndexMap::new(),
monitoring_address: None,
enable_span: false,
}
}
pub fn with_enable_span(mut self, enable: bool) -> Self {
self.enable_span = enable;
self
}
pub fn set_root_level(mut self, level: LevelFilter) -> Self {
self.root_level = level;
self
}
pub fn add_console_appender(mut self, appender_name: &str, out: Target, level: LevelFilter) -> Self {
self.targets.push(WriteTarget::Console(appender_name.to_string(), out, level));
self
}
pub fn add_file_appender(
mut self,
appender_name: &str,
file_name: &str,
level: LevelFilter,
max_file_size: usize,
max_file_history: usize,
) -> Self {
self.targets.push(WriteTarget::RollingFile(
appender_name.to_string(),
file_name.to_string(),
level,
MaxFileSize(max_file_size),
MaxFileHistory(max_file_history),
));
self
}
#[cfg(feature = "file-store-mongodb")]
pub async fn add_mongodb_appender(
mut self,
app_config: &crate::tina::server::application::AppConfig,
appender_name: &str,
collection_name: &str,
display_label: &str,
level: LevelFilter,
expires: std::time::Duration,
) -> Self {
let database = app_config.mongodb.build_log_db().await.expect("build log database failed");
crate::tina::log::mongodb::add_logger(collection_name, display_label);
self.targets.push(WriteTarget::Mongodb(appender_name.to_string(), collection_name.to_string(), database, level, expires));
self
}
pub fn set_root_appender(mut self, appender_names: Vec<&str>) -> Self {
self.root_config = appender_names.into_iter().map(|s| s.to_string()).collect();
self
}
pub fn set_module_appender(mut self, module: &str, appender_names: Vec<&str>, level: LevelFilter, additive: bool) -> Self {
self.custom_config.insert(
module.to_string(),
CustomLogConfig(module.to_string(), level, appender_names.into_iter().map(|v| v.to_string()).collect(), additive),
);
self
}
pub fn monitoring_address(mut self, monitoring_address: &str) -> Self {
let monitoring_address = SocketAddr::from_str(monitoring_address)
.unwrap_or_else(|err| panic!("parse address failed: {}, reason: {}", monitoring_address, err));
self.monitoring_address = Some(monitoring_address);
self
}
pub async fn init(self) -> AppResult<()> {
Log::init_logger_handle(self).await
}
}
impl Default for LogConfig {
fn default() -> Self {
LogConfig::new("./log", "default-server")
.add_console_appender("default", Target::Stdout, LevelFilter::Trace)
.set_root_appender(vec!["default"])
.set_root_level(LevelFilter::Debug)
}
}
pub struct Log;
const CHANNEL_CAPACITY: usize = 1024;
#[derive(Debug, Clone, Copy)]
enum LogPattern {
Simple,
SimpleWithSpan,
Advance,
AdvanceWithSpan,
}
impl AsRef<str> for LogPattern {
fn as_ref(&self) -> &str {
match self {
LogPattern::Simple => "{X(business_thread_time)} {h({l:<5})} [ {t:>30.30} - {h({L:<4})} ]{X(request-id)} - {h({m})}{n}",
LogPattern::SimpleWithSpan => "{X(business_thread_time)} {h({l:<5})} [ {t:>30.30} - {h({L:<4})} ]{X(span-data)} - {h({m})}{n}",
LogPattern::Advance => {
"{X(business_thread_time)} [ {X(business_thread_name)} ] {h({l})} [ {t} - {h({L})} ]{X(request-id)} - {h({m})}{n}"
}
LogPattern::AdvanceWithSpan => {
"{X(business_thread_time)} [ {X(business_thread_name)} ] {h({l})} [ {t} - {h({L})} ]{X(span-data)} - {h({m})}{n}"
}
}
}
}
fn get_log_pattern(config: &LogConfig) -> LogPattern {
match config.enable_span {
true => match cfg!(debug_assertions) {
true => LogPattern::SimpleWithSpan,
false => LogPattern::AdvanceWithSpan,
},
false => match cfg!(debug_assertions) {
true => LogPattern::Simple,
false => LogPattern::Advance,
},
}
}
impl Log {
pub async fn init_logger_handle(init_config: LogConfig) -> AppResult<()> {
let pattern = get_log_pattern(&init_config);
let log_home = init_config.log_home;
let web_server_name = init_config.web_server_name;
let targets = init_config.targets;
let root_config = init_config.root_config;
let root_level = init_config.root_level;
let custom_config = init_config.custom_config;
let mut config_builder = Box::new(Config::builder());
for config in targets.into_iter() {
match config {
WriteTarget::Console(appender_name, target, level_filter) => {
let appender =
ConsoleAppender::builder().encoder(Box::new(PatternEncoder::new(pattern.as_ref()))).target(target).build();
let appender = AsyncAppender::new(Arc::new(appender), CHANNEL_CAPACITY);
*config_builder = config_builder.appender(
Appender::builder()
.filter(Box::new(ThresholdFilter::new(level_filter)))
.build(appender_name.clone(), Box::new(appender)),
);
}
WriteTarget::RollingFile(appender_name, file_name, level_filter, max_file_size, max_file_history) => {
let appender = Log::new_rolling_file_appender(
log_home.as_str(),
web_server_name.as_str(),
&file_name,
max_file_size.0 as u64,
max_file_history.0 as u32,
pattern,
)
.expect("创建RollingFileAppender失败...");
*config_builder = config_builder.appender(
Appender::builder()
.filter(Box::new(ThresholdFilter::new(level_filter)))
.build(appender_name.clone(), Box::new(appender)),
);
}
#[cfg(feature = "file-store-mongodb")]
WriteTarget::Mongodb(appender_name, collection_name, database, level_filter, expires) => {
let appender = Log::new_mongodb_appender(collection_name.as_str(), database, level_filter, expires)
.await
.expect("创建MongodbAsyncAppender失败");
*config_builder = config_builder.appender(
Appender::builder().filter(Box::new(ThresholdFilter::new(level_filter))).build(appender_name, Box::new(appender)),
);
}
};
}
let root_level = Rc::new(root_level);
let mut root_builder = Box::new(Root::builder());
for appender_name in root_config.into_iter() {
*root_builder = root_builder.appender(appender_name);
}
for (_, log) in custom_config.into_iter() {
let mut logger_builder = Box::new(Logger::builder());
let name = log.0;
let level = Rc::new(log.1);
let additive = log.3;
for appender_name in log.2.into_iter() {
*logger_builder = logger_builder.appender(appender_name).additive(additive);
}
let logger = logger_builder.build(name, *level.as_ref());
*config_builder = config_builder.logger(logger);
}
let log4s_config = config_builder.build(root_builder.build(*root_level.as_ref())).expect("初始化log4rs配置失败!");
let logger = log4rs::Logger::new(log4s_config);
let max_level = logger.max_log_level();
let logger = AsyncLogger::new(logger, max_level.to_level().ok_or_else(|| app_system_error!("No Log Level"))?, CHANNEL_CAPACITY);
log::set_max_level(max_level);
log::set_boxed_logger(Box::new(logger.clone())).map_err(app_error_from!())?;
let registry = tracing_subscriber::registry().with(logger);
tracing::subscriber::set_global_default(registry).map_err(app_error_from!())?;
Ok(())
}
fn new_rolling_file_appender(
log_home: &str,
web_server_name: &str,
file_name: &str,
trigger_size_byte: u64,
rolling_file_count: u32,
pattern: LogPattern,
) -> io::Result<AsyncAppender> {
let trigger = Box::new(SizeTrigger::new(trigger_size_byte)) as Box<dyn Trigger>;
let roller = Box::new(
FixedWindowRoller::builder()
.build((String::from(log_home) + "/" + web_server_name + "/gz/" + file_name + "-{}.gz").as_str(), rolling_file_count)
.expect("创建FixedWindowRoller失败..."),
);
let policy = CompoundPolicy::new(trigger, roller);
let appender = RollingFileAppender::builder()
.encoder(Box::new(PatternEncoder::new(pattern.as_ref())))
.append(true)
.build(String::from(log_home) + "/" + web_server_name + "/" + file_name, Box::new(policy))?;
let appender = Arc::new(appender);
let mut lock = GLOBAL_CLOSEABLE_APPENDERS.lock().map_err(|err| io::Error::new(io::ErrorKind::Other, err.to_string()))?;
lock.push(appender.clone() as Arc<dyn CloseableAppender>);
Ok(AsyncAppender::new(appender, CHANNEL_CAPACITY))
}
#[cfg(feature = "file-store-mongodb")]
async fn new_mongodb_appender(
appender_name: &str,
database: ::mongodb::Database,
level: LevelFilter,
expires: std::time::Duration,
) -> io::Result<AsyncAppender> {
let appender = mongodb::MongodbAppender::new(appender_name, database, level, expires).await.expect("创建MongodbAppender失败");
let appender = Arc::new(appender);
let mut lock = GLOBAL_CLOSEABLE_APPENDERS.lock().map_err(|err| io::Error::new(io::ErrorKind::Other, err.to_string()))?;
lock.push(appender.clone() as Arc<dyn CloseableAppender>);
Ok(AsyncAppender::new(appender, CHANNEL_CAPACITY))
}
}
#[derive(Clone, Debug)]
enum Command {
Record(AsyncRecord),
Flush,
Exit,
}
#[derive(Clone, Debug)]
struct AsyncRecord {
level: log::Level,
target: String,
args: String,
module_path: Option<String>,
file: Option<String>,
line: Option<u32>,
mdc: IndexMap<String, String>,
}
impl AsyncRecord {
fn from_record(record: &Record) -> AsyncRecord {
let mut mdc: IndexMap<String, String> = IndexMap::new();
log_mdc::iter(|key, value| {
mdc.insert(key.to_owned(), value.to_owned());
});
AsyncRecord {
level: record.level(),
target: record.metadata().target().to_string(),
args: record.args().to_string(),
module_path: record.module_path().map(|s| s.to_owned()),
file: record.file().map(|s| s.to_owned()),
line: record.line(),
mdc,
}
}
fn from_event(event: &tracing::Event<'_>) -> AsyncRecord {
let mut mdc: IndexMap<String, String> = IndexMap::new();
log_mdc::iter(|key, value| {
mdc.insert(key.to_owned(), value.to_owned());
});
let mut args = StringFields(None);
event.record(&mut args);
let metadata = event.metadata();
let level = match *metadata.level() {
tracing::Level::TRACE => Level::Trace,
tracing::Level::DEBUG => Level::Debug,
tracing::Level::INFO => Level::Info,
tracing::Level::WARN => Level::Warn,
tracing::Level::ERROR => Level::Error,
};
AsyncRecord {
level,
target: metadata.target().to_string(),
args: args.0.expect("no args found"),
module_path: metadata.module_path().map(|s| s.to_owned()),
file: metadata.file().map(|s| s.to_owned()),
line: metadata.line(),
mdc,
}
}
}
#[derive(Debug, Clone)]
pub struct AsyncLogger {
log_level: Level,
tracing_level: tracing::Level,
sender: Arc<SyncSender<Command>>,
span_data: Arc<RwLock<Option<Arc<RwLock<AsyncLoggerFields>>>>>,
}
#[derive(Debug, Clone)]
pub struct AsyncLoggerFields(Vec<(Arc<String>, Arc<String>)>);
#[derive(Clone, prost::Message)]
pub struct AsyncLoggerFieldsMessage {
#[prost(message, repeated, tag = "1")]
pub fields: Vec<AsyncLoggerFieldMessage>,
}
#[derive(Clone, prost::Message)]
pub struct AsyncLoggerFieldMessage {
#[prost(string, tag = "1")]
pub key: prost::alloc::string::String,
#[prost(string, tag = "2")]
pub value: prost::alloc::string::String,
}
#[derive(Debug)]
struct StringFields(Option<String>);
#[derive(Debug)]
struct AsyncAppender {
sender: SyncSender<Command>,
}
impl AsyncLogger {
fn new(logger: log4rs::Logger, level: Level, capacity: usize) -> AsyncLogger {
let tracing_level = match level {
Level::Error => tracing::Level::ERROR,
Level::Warn => tracing::Level::WARN,
Level::Info => tracing::Level::INFO,
Level::Debug => tracing::Level::DEBUG,
Level::Trace => tracing::Level::TRACE,
};
let (sender, receiver) = std::sync::mpsc::sync_channel(capacity);
let receiver = receiver;
let logger = Arc::new(logger);
let async_logger = AsyncLogger {
log_level: level,
sender: Arc::new(sender),
tracing_level,
span_data: Arc::new(RwLock::new(None)),
};
std::thread::Builder::new()
.name("日志线程".to_owned())
.spawn(move || loop {
if Utility::is_terminated() {
return;
}
match receiver.recv() {
Ok(cmd) => match cmd {
Command::Record(async_record) => {
let AsyncRecord {
level,
target,
args,
module_path,
file,
line,
mdc,
} = async_record;
log_mdc::clear();
for (key, value) in mdc.into_iter() {
log_mdc::insert(key, value);
}
log::Log::log(
logger.as_ref(),
&Record::builder()
.level(level)
.target(target.as_str())
.args(format_args!("{}", args))
.module_path(module_path.as_deref())
.file(file.as_deref())
.line(line)
.build(),
)
}
Command::Flush => {
log::Log::flush(logger.as_ref());
}
Command::Exit => {
break;
}
},
Err(err) => {
println!("receive log data failed: {}", err);
}
}
})
.expect("启动日志线程失败!");
async_logger
}
}
impl AsyncAppender {
fn new(appender: Arc<dyn Append>, capacity: usize) -> AsyncAppender {
let (sender, receiver) = std::sync::mpsc::sync_channel(capacity);
let receiver = receiver;
let async_appender = AsyncAppender {
sender,
};
std::thread::Builder::new()
.name("日志Append线程".to_owned())
.spawn(move || loop {
if Utility::is_terminated() {
return;
}
match receiver.recv() {
Ok(cmd) => match cmd {
Command::Record(async_record) => {
let AsyncRecord {
level,
target,
args,
module_path,
file,
line,
mdc,
} = async_record;
log_mdc::clear();
for (key, value) in mdc.into_iter() {
log_mdc::insert(key, value);
}
match appender.append(
&Record::builder()
.level(level)
.target(target.as_str())
.args(format_args!("{}", args))
.module_path(module_path.as_deref())
.file(file.as_deref())
.line(line)
.build(),
) {
Ok(_) => {}
Err(err) => {
eprintln!("{}", err);
}
}
}
Command::Flush => {
appender.flush();
}
Command::Exit => {
break;
}
},
Err(err) => {
println!("receive log data failed: {}", err);
}
}
})
.expect("启动日志Append线程失败!");
async_appender
}
}
impl log::Log for AsyncLogger {
fn enabled(&self, metadata: &Metadata) -> bool {
metadata.level() <= self.log_level
}
fn log(&self, record: &Record) {
if log::Log::enabled(&self, record.metadata()) {
let level = record.level();
if self.log_level < level {
return;
}
log_mdc::insert("business_thread_id", thread_id::get().to_string());
log_mdc::insert("business_thread_name", current().name().unwrap_or("unkown").to_string());
log_mdc::insert("business_thread_time", Local::now().format("%Y-%m-%d %H:%M:%S%.3f").to_string());
{
let span = self.span_data.read().expect("lock span data failed");
if let Some(fields) = span.as_ref() {
let fields = fields.read().expect("lock span fields failed");
log_mdc::insert(Constants::SPAN_DATA_HEADER, fields.clone());
for (k, v) in fields.0.iter() {
let k = k.as_str();
let v = v.as_str();
if k == "request_id" {
log_mdc::insert(Constants::REQUEST_ID_HEADER, format!(" [ {} ]", v));
log_mdc::insert(Constants::REQUEST_ID_HEADER_VALUE, v);
} else {
log_mdc::insert(k, v);
}
}
}
}
let async_record = AsyncRecord::from_record(record);
match self.sender.send(Command::Record(async_record)) {
Ok(_) => {}
Err(err) => {
println!("AsyncLogger log failed: {}", err);
}
}
}
}
fn flush(&self) {
match self.sender.send(Command::Flush) {
Ok(_) => {}
Err(err) => {
println!("AsyncLogger flush failed: {}", err);
}
}
}
}
impl AsyncLoggerFields {
pub fn lookup_field_values(&self, field: &str) -> Vec<&str> {
let mut datas = Vec::new();
for item in self.0.iter() {
if item.0.as_str() == field {
datas.push(item.1.as_str());
}
}
datas
}
pub fn merge_to_end(&mut self, other: &AsyncLoggerFields) {
self.0.extend_from_slice(other.0.as_slice());
}
pub fn merge_to_start(&mut self, other: &AsyncLoggerFields) {
let mut datas = Vec::with_capacity(self.0.len() + other.0.len());
let mut cache = Vec::new();
std::mem::swap(&mut self.0, &mut cache);
datas.extend_from_slice(other.0.as_slice());
datas.append(&mut cache);
self.0 = datas;
}
pub fn merge_into_start(&mut self, other: AsyncLoggerFields) {
let mut data = other.0;
let mut cache = Vec::new();
std::mem::swap(&mut self.0, &mut cache);
data.append(&mut cache);
self.0 = data;
}
pub fn into_message(self) -> AsyncLoggerFieldsMessage {
AsyncLoggerFieldsMessage::from(self)
}
pub fn from_message(message: AsyncLoggerFieldsMessage) -> Self {
Self::from(message)
}
}
impl From<AsyncLoggerFields> for String {
fn from(value: AsyncLoggerFields) -> Self {
let mut sj = StringJoiner::new(Some("[ "), Some(" ]"), Some(", "));
for (k, v) in value.0.into_iter() {
sj.push_string(format!("{k} = {v}"));
}
sj.into_string()
}
}
impl From<AsyncLoggerFields> for AsyncLoggerFieldsMessage {
fn from(value: AsyncLoggerFields) -> Self {
let fields = value
.0
.into_iter()
.map(|v| AsyncLoggerFieldMessage {
key: v.0.to_string(),
value: v.1.to_string(),
})
.collect::<Vec<AsyncLoggerFieldMessage>>();
Self {
fields,
}
}
}
impl From<AsyncLoggerFieldsMessage> for AsyncLoggerFields {
fn from(value: AsyncLoggerFieldsMessage) -> Self {
let fields = value.fields.into_iter().map(|v| (Arc::new(v.key), Arc::new(v.value))).collect::<Vec<(Arc<String>, Arc<String>)>>();
Self(fields)
}
}
impl tracing::field::Visit for AsyncLoggerFields {
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
let name = field.as_ref();
self.0.push((Arc::new(name.to_owned()), Arc::new(format!("{:?}", value))));
}
fn record_str(&mut self, field: &tracing::field::Field, value: &str) {
let name = field.as_ref();
self.0.push((Arc::new(name.to_owned()), Arc::new(value.to_owned())));
}
}
impl tracing::field::Visit for StringFields {
fn record_debug(&mut self, _field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
self.0 = Some(format!("{:?}", value));
}
fn record_str(&mut self, _field: &tracing::field::Field, value: &str) {
self.0 = Some(value.to_owned());
}
}
impl<S> tracing_subscriber::Layer<S> for AsyncLogger
where
S: tracing::Subscriber + for<'lookup> tracing_subscriber::registry::LookupSpan<'lookup>,
{
fn enabled(&self, metadata: &tracing::Metadata<'_>, _ctx: tracing_subscriber::layer::Context<'_, S>) -> bool {
metadata.level() <= &self.tracing_level
}
fn on_new_span(&self, attrs: &tracing::span::Attributes<'_>, id: &tracing::span::Id, ctx: tracing_subscriber::layer::Context<'_, S>) {
let span = ctx.span(id).unwrap_or_else(|| panic!("no span found: {:?}", id));
let mut fields = AsyncLoggerFields(vec![]);
if let Some(parent) = span.parent() {
let ext = parent.extensions();
if let Some(parent_fields) = ext.get::<Arc<RwLock<AsyncLoggerFields>>>() {
let parent_fields = parent_fields.clone();
let parent_fields = parent_fields.read();
if let Ok(parent_fields) = parent_fields {
fields.merge_to_end(&parent_fields);
}
}
}
attrs.record(&mut fields);
span.extensions_mut().insert(Arc::new(RwLock::new(fields)));
}
fn on_enter(&self, id: &tracing::span::Id, ctx: tracing_subscriber::layer::Context<'_, S>) {
let span_ref = ctx.span(id).unwrap_or_else(|| panic!("no span found: {:?}", id));
let ext = span_ref.extensions();
let fields = ext.get::<Arc<RwLock<AsyncLoggerFields>>>().unwrap_or_else(|| panic!("no AsyncLoggerFields in span: {:?}", id));
let mut span_data = self.span_data.write().expect("lock span data failed");
(*span_data) = Some(fields.clone());
}
fn on_exit(&self, _id: &tracing::span::Id, _ctx: tracing_subscriber::layer::Context<'_, S>) {
let mut span_data = self.span_data.write().expect("lock span data failed");
(*span_data) = None;
}
fn on_record(&self, span: &tracing::span::Id, values: &tracing::span::Record<'_>, ctx: tracing_subscriber::layer::Context<'_, S>) {
let span_ref = ctx.span(span).unwrap_or_else(|| panic!("no span found: {:?}", span));
let ext = span_ref.extensions();
let fields = ext.get::<Arc<RwLock<AsyncLoggerFields>>>().unwrap_or_else(|| panic!("no AsyncLoggerFields in span: {:?}", span));
let mut fields = fields.write().expect("lock span fields failed");
values.record(&mut *fields);
}
fn on_event(&self, event: &tracing::Event<'_>, ctx: tracing_subscriber::layer::Context<'_, S>) {
if tracing_subscriber::Layer::enabled(self, event.metadata(), ctx) {
log_mdc::insert("business_thread_id", thread_id::get().to_string());
log_mdc::insert("business_thread_name", current().name().unwrap_or("unkown").to_string());
log_mdc::insert("business_thread_time", Local::now().format("%Y-%m-%d %H:%M:%S%.3f").to_string());
{
let span = self.span_data.read().expect("lock span data failed");
if let Some(fields) = span.as_ref() {
let fields = fields.read().expect("lock span fields failed");
log_mdc::insert(Constants::SPAN_DATA_HEADER, fields.clone());
for (k, v) in fields.0.iter() {
let k = k.as_str();
let v = v.as_str();
if k == "request_id" {
log_mdc::insert(Constants::REQUEST_ID_HEADER, format!(" [ {} ]", v));
log_mdc::insert(Constants::REQUEST_ID_HEADER_VALUE, v);
} else {
log_mdc::insert(k, v);
}
}
}
}
let async_record = AsyncRecord::from_event(event);
match self.sender.send(Command::Record(async_record)) {
Ok(_) => {}
Err(err) => {
println!("AsyncLogger log failed: {}", err);
}
}
}
}
fn on_follows_from(
&self,
_span: &tracing_core::span::Id,
follows: &tracing_core::span::Id,
ctx: tracing_subscriber::layer::Context<'_, S>,
) {
if let Some(span_ref) = ctx.span(follows) {
let ext = span_ref.extensions();
if let Some(fields) = ext.get::<Arc<RwLock<AsyncLoggerFields>>>() {
let follows_span_data = {
let lock = fields.read().expect("lock follows span data failed");
lock.clone()
};
let mut span_data = self.span_data.write().expect("lock span data failed");
match span_data.as_mut() {
Some(span_data) => {
let mut lock = span_data.write().expect("lock span inner data failed");
lock.merge_into_start(follows_span_data);
}
None => {
(*span_data) = Some(Arc::new(RwLock::new(follows_span_data)));
}
}
}
}
}
}
impl Append for AsyncAppender {
fn append(&self, record: &Record) -> anyhow::Result<()> {
log_mdc::insert("business_thread_id", thread_id::get().to_string());
log_mdc::insert("business_thread_time", Local::now().format("%Y-%m-%d %H:%M:%S%.3f").to_string());
let async_record = AsyncRecord::from_record(record);
match self.sender.send(Command::Record(async_record)) {
Ok(_) => {}
Err(err) => {
println!("AsyncAppender log failed: {}", err);
}
}
Ok(())
}
fn flush(&self) {
match self.sender.send(Command::Flush) {
Ok(_) => {}
Err(err) => {
println!("AsyncAppender flush failed: {}", err);
}
}
}
}
impl Drop for AsyncLogger {
fn drop(&mut self) {
match self.sender.send(Command::Exit) {
Ok(_) => {}
Err(err) => {
println!("AsyncLogger drop failed: {}", err);
}
}
}
}
impl Drop for AsyncAppender {
fn drop(&mut self) {
match self.sender.send(Command::Exit) {
Ok(_) => {}
Err(err) => {
println!("AsyncAppender drop failed: {}", err);
}
}
}
}
impl CloseableAppender for RollingFileAppender {
fn shutdown(&self) {
println!("停止RollingFileAppender开始: {:?}", self);
self.flush();
println!("停止RollingFileAppender完成: {:?}", self);
}
}
#[cfg(test)]
mod tests {
#![allow(unused)]
use tokio::runtime::Runtime;
use tracing::{subscriber::NoSubscriber, Id};
use super::*;
fn init() -> Runtime {
let runtime = tokio::runtime::Builder::new_multi_thread().worker_threads(20).enable_all().build().expect("build runtime failed");
runtime.block_on(async move {
Log::init_logger_handle(
LogConfig::new(".", "abc")
.add_console_appender("stdout", Target::Stdout, LevelFilter::Trace)
.set_root_appender(vec!["stdout"]),
)
.await
.expect("init log failed");
});
runtime
}
#[test]
fn test_span() {
init();
let span = tracing::debug_span!("test", foo = "abc");
span.record("foo", "abcd");
span.in_scope(|| {
let span = tracing::Span::current();
tracing::info!("hello, {:?}", span);
})
}
}