1use std::{
8 fmt::Write as _,
9 future::Future,
10 io::{self, stderr, stdout, Write},
11 pin::Pin,
12 str::FromStr,
13 thread::{self, JoinHandle},
14};
15
16use crate::Result;
17use chrono::{DateTime, Local, Utc};
18use colored::{Color, Colorize as _};
19use futures::executor::block_on;
20use serde::{Deserialize, Serialize};
21use serde_json::{Map, Value};
22use serde_with::{serde_as, DisplayFromStr};
23use tokio::{
24 select,
25 sync::mpsc::{unbounded_channel, UnboundedSender},
26};
27use tracing::Level;
28
29#[serde_as]
31#[derive(Serialize, Deserialize, Debug)]
32pub struct LogItem {
33 pub time: Value,
36 #[serde_as(as = "DisplayFromStr")]
37 pub level: Level,
38 pub message: String,
39 #[serde(skip_serializing_if = "String::is_empty")]
40 pub target: String,
41 #[serde(skip_serializing_if = "Map::is_empty")]
43 pub fields: Map<String, Value>,
44 #[serde(skip_serializing_if = "Map::is_empty")]
46 pub span: Map<String, Value>,
47 #[serde(skip_serializing_if = "Option::is_none")]
48 pub filename: Option<String>,
49 #[serde(skip_serializing_if = "Option::is_none")]
50 pub line_number: Option<i64>,
51}
52
53impl LogItem {
54 fn json_take_object(mp: &mut Map<String, Value>, key: &str) -> Map<String, Value> {
55 if let Value::Object(x) = mp.remove(key).unwrap_or_default() {
56 x
57 } else {
58 Map::default()
59 }
60 }
61
62 fn from_json(mut s: Map<String, Value>) -> Self {
63 let target = s
64 .get("target")
65 .and_then(Value::as_str)
66 .unwrap_or_default()
67 .to_owned();
68 let level = Level::from_str(s.get("level").and_then(Value::as_str).unwrap_or("ERROR"))
69 .unwrap_or(Level::ERROR);
70 let filename = s
71 .get("filename")
72 .and_then(Value::as_str)
73 .map(str::to_string);
74 let line_number = s.get("line_number").and_then(Value::as_i64);
75 let mut fields = Self::json_take_object(&mut s, "fields");
76 let message = fields
77 .remove("message")
78 .unwrap_or_default()
79 .as_str()
80 .unwrap_or_default()
81 .to_owned();
82 let span = Self::json_take_object(&mut s, "span");
83 Self {
84 time: Value::default(),
85 level,
86 message,
87 target,
88 fields,
89 span,
90 filename,
91 line_number,
92 }
93 }
94}
95
96struct LogSender {
97 tx: UnboundedSender<Map<String, Value>>,
98}
99
100impl Write for LogSender {
101 fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
102 self.tx
103 .send(serde_json::from_slice(buf)?)
104 .or(Err(io::ErrorKind::BrokenPipe))?;
105 Ok(buf.len())
106 }
107
108 fn flush(&mut self) -> io::Result<()> {
109 Ok(())
111 }
112}
113
114impl LogSender {
115 fn new(tx: UnboundedSender<Map<String, Value>>) -> impl Fn() -> Self {
116 move || Self { tx: tx.clone() }
117 }
118}
119
120fn init_subscriber(
122 level: Level,
123 filename: bool,
124 line_number: bool,
125 tx: UnboundedSender<Map<String, Value>>,
126) {
127 tracing_subscriber::fmt()
128 .with_max_level(level)
129 .with_writer(LogSender::new(tx))
130 .without_time()
131 .with_file(filename)
132 .with_line_number(line_number)
133 .json()
134 .init();
135}
136
137#[derive(Clone)]
140pub struct Logger {
141 tx: UnboundedSender<Map<String, Value>>,
142}
143
144impl Logger {
145 pub fn sender(&self) -> UnboundedSender<Map<String, Value>> {
147 self.tx.clone()
148 }
149
150 pub fn init(&self, builder: &LoggerBuilder) {
153 init_subscriber(
154 builder.level,
155 builder.filename,
156 builder.line_number,
157 self.tx.clone(),
158 );
159 }
160}
161
162pub type WriterFn = Box<dyn Fn(LogItem, Box<dyn Write>) -> Result<()> + Send>;
164pub type FilterFn = Box<dyn Fn(&LogItem) -> bool + Send>;
166pub type TransformerFn = Box<dyn Fn(LogItem) -> LogItem + Send>;
168pub type HandlerFn = Box<dyn Fn(&Map<String, Value>) -> Pin<Box<dyn Future<Output = bool>>> + Send>;
171
172pub struct LoggerGuard {
178 stop_tx: UnboundedSender<()>,
179 join: Option<JoinHandle<()>>,
180}
181
182impl Drop for LoggerGuard {
183 fn drop(&mut self) {
184 let _ = self.stop_tx.send(());
185 if let Some(x) = self.join.take() {
186 let _ = x.join();
187 }
188 }
189}
190
191pub struct LoggerBuilder {
192 json: bool,
193 level: Level,
194 filename: bool,
195 line_number: bool,
196 filter: Option<FilterFn>,
197 transformer: Option<TransformerFn>,
198 json_writer: WriterFn,
199 color_writer: WriterFn,
200 handler: Option<HandlerFn>,
201}
202
203impl Default for LoggerBuilder {
204 fn default() -> Self {
205 Self::new()
206 }
207}
208
209impl LoggerBuilder {
210 pub fn fmt_level(level: &Level) -> String {
217 format!("{: >5}", level.to_string())
218 .bold()
219 .color(match *level {
220 Level::TRACE | Level::DEBUG => Color::Magenta,
221 Level::INFO => Color::Green,
222 Level::WARN => Color::Yellow,
223 Level::ERROR => Color::Red,
224 })
225 .to_string()
226 }
227
228 fn default_json_writer(item: LogItem, mut writer: Box<dyn Write>) -> Result<()> {
229 let v = serde_json::to_string(&item).unwrap_or_default();
230 writer.write_fmt(format_args!("{v}\n"))?;
231 writer.flush().map_err(Into::into)
232 }
233
234 fn default_color_writer(item: LogItem, mut writer: Box<dyn Write>) -> Result<()> {
235 let mut buf = String::new();
236 write!(
237 buf,
238 "{} {} {}",
239 item.time.as_str().unwrap_or_default().bright_black(),
240 Self::fmt_level(&item.level),
241 item.target.bright_black()
242 )?;
243 if let (Some(filename), Some(line_number)) = (item.filename, item.line_number) {
244 buf += &format!("({filename}:{line_number})")
245 .bright_black()
246 .to_string();
247 }
248 write!(buf, "{} {}", ":".bright_black(), item.message)?;
249 for (k, v) in &item.fields {
250 if !k.starts_with("log.") {
251 buf += &format!(" field.{k}={v}").bright_black().to_string();
252 }
253 }
254 for (k, v) in item.span {
255 if !k.starts_with("http.") && !k.starts_with("otel.") && k != "name" {
256 buf += &format!(" span.{k}={v}").bright_black().to_string();
257 }
258 }
259
260 writer.write_fmt(format_args!("{buf}\n"))?;
261 writer.flush().map_err(Into::into)
262 }
263
264 pub fn new() -> Self {
267 Self {
268 json: false,
269 level: Level::INFO,
270 filename: false,
271 line_number: false,
272 filter: None,
273 transformer: None,
274 json_writer: Box::new(Self::default_json_writer),
275 color_writer: Box::new(Self::default_color_writer),
276 handler: None,
277 }
278 }
279
280 pub fn json_writer(mut self, writer: WriterFn) -> Self {
285 self.json_writer = writer;
286 self
287 }
288
289 pub fn color_writer(mut self, writer: WriterFn) -> Self {
294 self.color_writer = writer;
295 self
296 }
297
298 pub fn json(mut self) -> Self {
300 self.json = true;
301 self
302 }
303
304 pub fn level(mut self, level: Level) -> Self {
306 self.level = level;
307 self
308 }
309
310 pub fn filename(mut self) -> Self {
312 self.filename = true;
313 self
314 }
315
316 pub fn line_number(mut self) -> Self {
318 self.line_number = true;
319 self
320 }
321
322 pub fn handler<F>(mut self, handler: F) -> Self
331 where
332 F: Fn(&Map<String, Value>) -> Pin<Box<dyn Future<Output = bool>>> + Send + 'static,
333 {
334 self.handler = Some(Box::new(handler));
335 self
336 }
337
338 pub fn filter<F>(mut self, filter: F) -> Self
345 where
346 F: Fn(&LogItem) -> bool + Send + 'static,
347 {
348 self.filter = Some(Box::new(filter));
349 self
350 }
351
352 pub fn transformer<F>(mut self, transformer: F) -> Self
359 where
360 F: Fn(LogItem) -> LogItem + Send + 'static,
361 {
362 self.transformer = Some(Box::new(transformer));
363 self
364 }
365
366 pub fn start(self) -> (Logger, LoggerGuard) {
373 let (tx, mut rx) = unbounded_channel();
374 let (stop_tx, mut stop_rx) = unbounded_channel();
375 init_subscriber(self.level, self.filename, self.line_number, tx.clone());
376
377 let join = thread::spawn(move || {
378 let handler = |v: Map<String, Value>| async {
379 if let Some(x) = &self.handler {
380 if !x(&v).await {
381 return;
382 }
383 }
384 let mut item = LogItem::from_json(v);
385 let time = item.fields.remove("_time").unwrap_or_default().as_i64();
386 if self.json {
387 item.time = time.unwrap_or_else(|| Utc::now().timestamp_micros()).into();
388 } else {
389 item.time = time
390 .map_or_else(Local::now, |v| {
391 DateTime::from_timestamp_micros(v)
392 .unwrap_or_default()
393 .into()
394 })
395 .format("%F %T%.6f")
396 .to_string()
397 .into();
398 }
399
400 if let Some(filter) = &self.filter {
401 if !filter(&item) {
402 return;
403 }
404 }
405 if let Some(transformer) = &self.transformer {
406 item = transformer(item);
407 }
408 let writer: Box<dyn io::Write> = if item.level <= Level::WARN {
409 Box::new(stderr())
410 } else {
411 Box::new(stdout())
412 };
413 if self.json {
414 let _ = (self.json_writer)(item, writer);
415 } else {
416 let _ = (self.color_writer)(item, writer);
417 }
418 };
419 block_on(async move {
420 loop {
421 select! {
422 Some(v) = rx.recv() => {
423 handler(v).await;
424 },
425 _ = stop_rx.recv() => {
426 while let Ok(v) = rx.try_recv(){
427 handler(v).await;
428 }
429 break;
430 }
431 }
432 }
433 })
434 });
435 (
436 Logger { tx },
437 LoggerGuard {
438 stop_tx,
439 join: Some(join),
440 },
441 )
442 }
443}