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]
30#[derive(Serialize, Deserialize, Debug)]
31pub struct LogItem {
32 pub time: Value,
33 #[serde_as(as = "DisplayFromStr")]
34 pub level: Level,
35 pub message: String,
36 #[serde(skip_serializing_if = "String::is_empty")]
37 pub target: String,
38 #[serde(skip_serializing_if = "Map::is_empty")]
39 pub fields: Map<String, Value>,
40 #[serde(skip_serializing_if = "Map::is_empty")]
41 pub span: Map<String, Value>,
42 #[serde(skip_serializing_if = "Option::is_none")]
43 pub filename: Option<String>,
44 #[serde(skip_serializing_if = "Option::is_none")]
45 pub line_number: Option<i64>,
46}
47
48impl LogItem {
49 fn json_take_object(mp: &mut Map<String, Value>, key: &str) -> Map<String, Value> {
50 if let Value::Object(x) = mp.remove(key).unwrap_or_default() {
51 x
52 } else {
53 Map::default()
54 }
55 }
56
57 fn from_json(mut s: Map<String, Value>) -> Self {
58 let target = s
59 .get("target")
60 .and_then(Value::as_str)
61 .unwrap_or_default()
62 .to_owned();
63 let level = Level::from_str(s.get("level").and_then(Value::as_str).unwrap_or("ERROR"))
64 .unwrap_or(Level::ERROR);
65 let filename = s
66 .get("filename")
67 .and_then(Value::as_str)
68 .map(str::to_string);
69 let line_number = s.get("line_number").and_then(Value::as_i64);
70 let mut fields = Self::json_take_object(&mut s, "fields");
71 let message = fields
72 .remove("message")
73 .unwrap_or_default()
74 .as_str()
75 .unwrap_or_default()
76 .to_owned();
77 let span = Self::json_take_object(&mut s, "span");
78 Self {
79 time: Value::default(),
80 level,
81 message,
82 target,
83 fields,
84 span,
85 filename,
86 line_number,
87 }
88 }
89}
90
91struct LogSender {
92 tx: UnboundedSender<Map<String, Value>>,
93}
94
95impl Write for LogSender {
96 fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
97 self.tx
98 .send(serde_json::from_slice(buf)?)
99 .or(Err(io::ErrorKind::BrokenPipe))?;
100 Ok(buf.len())
101 }
102
103 fn flush(&mut self) -> io::Result<()> {
104 Ok(())
106 }
107}
108
109impl LogSender {
110 fn new(tx: UnboundedSender<Map<String, Value>>) -> impl Fn() -> Self {
111 move || Self { tx: tx.clone() }
112 }
113}
114
115#[derive(Clone)]
116pub struct Logger {
117 tx: UnboundedSender<Map<String, Value>>,
118}
119
120impl Logger {
121 pub fn sender(&self) -> UnboundedSender<Map<String, Value>> {
123 self.tx.clone()
124 }
125
126 pub fn init(&self, builder: &LoggerBuilder) {
129 tracing_subscriber::fmt()
130 .with_max_level(builder.level)
131 .with_writer(LogSender::new(self.tx.clone()))
132 .without_time()
133 .with_file(builder.filename)
134 .with_line_number(builder.line_number)
135 .json()
136 .init();
137 }
138}
139
140pub type WriterFn = Box<dyn Fn(LogItem, Box<dyn Write>) -> Result<()> + Send>;
141pub type FilterFn = Box<dyn Fn(&LogItem) -> bool + Send>;
142pub type TransformerFn = Box<dyn Fn(LogItem) -> LogItem + Send>;
143pub type HandlerFn = Box<dyn Fn(&Map<String, Value>) -> Pin<Box<dyn Future<Output = bool>>> + Send>;
144
145pub struct LoggerGuard {
151 stop_tx: UnboundedSender<()>,
152 join: Option<JoinHandle<()>>,
153}
154
155impl Drop for LoggerGuard {
156 fn drop(&mut self) {
157 let _ = self.stop_tx.send(());
158 if let Some(x) = self.join.take() {
159 let _ = x.join();
160 }
161 }
162}
163
164pub struct LoggerBuilder {
165 json: bool,
166 level: Level,
167 filename: bool,
168 line_number: bool,
169 filter: Option<FilterFn>,
170 transformer: Option<TransformerFn>,
171 json_writer: WriterFn,
172 color_writer: WriterFn,
173 handler: Option<HandlerFn>,
174}
175
176impl Default for LoggerBuilder {
177 fn default() -> Self {
178 Self::new()
179 }
180}
181
182impl LoggerBuilder {
183 pub fn fmt_level(level: &Level) -> String {
190 format!("{: >5}", level.to_string())
191 .bold()
192 .color(match *level {
193 Level::TRACE | Level::DEBUG => Color::Magenta,
194 Level::INFO => Color::Green,
195 Level::WARN => Color::Yellow,
196 Level::ERROR => Color::Red,
197 })
198 .to_string()
199 }
200
201 fn default_json_writer(item: LogItem, mut writer: Box<dyn Write>) -> Result<()> {
202 let v = serde_json::to_string(&item).unwrap_or_default();
203 writer.write_fmt(format_args!("{v}\n"))?;
204 writer.flush().map_err(Into::into)
205 }
206
207 fn default_color_writer(item: LogItem, mut writer: Box<dyn Write>) -> Result<()> {
208 let mut buf = String::new();
209 write!(
210 buf,
211 "{} {} {}",
212 item.time.as_str().unwrap_or_default().bright_black(),
213 Self::fmt_level(&item.level),
214 item.target.bright_black()
215 )?;
216 if let Some(filename) = item.filename {
217 if let Some(line_number) = item.line_number {
218 buf += &format!("({}:{})", filename, line_number)
219 .bright_black()
220 .to_string();
221 }
222 }
223 write!(buf, "{} {}", ":".bright_black(), item.message)?;
224 for (k, v) in &item.fields {
225 if !k.starts_with("log.") {
226 buf += &format!(" field.{k}={v}").bright_black().to_string();
227 }
228 }
229 for (k, v) in item.span {
230 if !k.starts_with("http.") && !k.starts_with("otel.") && k != "name" {
231 buf += &format!(" span.{k}={v}").bright_black().to_string();
232 }
233 }
234
235 writer.write_fmt(format_args!("{buf}\n"))?;
236 writer.flush().map_err(Into::into)
237 }
238
239 pub fn new() -> Self {
242 Self {
243 json: false,
244 level: Level::INFO,
245 filename: false,
246 line_number: false,
247 filter: None,
248 transformer: None,
249 json_writer: Box::new(Self::default_json_writer),
250 color_writer: Box::new(Self::default_color_writer),
251 handler: None,
252 }
253 }
254
255 pub fn json_writer(mut self, writer: WriterFn) -> Self {
260 self.json_writer = writer;
261 self
262 }
263
264 pub fn color_writer(mut self, writer: WriterFn) -> Self {
269 self.color_writer = writer;
270 self
271 }
272
273 pub fn json(mut self) -> Self {
275 self.json = true;
276 self
277 }
278
279 pub fn level(mut self, level: Level) -> Self {
281 self.level = level;
282 self
283 }
284
285 pub fn filename(mut self) -> Self {
287 self.filename = true;
288 self
289 }
290
291 pub fn line_number(mut self) -> Self {
293 self.line_number = true;
294 self
295 }
296
297 pub fn handler<F>(mut self, handler: F) -> Self
306 where
307 F: Fn(&Map<String, Value>) -> Pin<Box<dyn Future<Output = bool>>> + Send + 'static,
308 {
309 self.handler = Some(Box::new(handler));
310 self
311 }
312
313 pub fn filter<F>(mut self, filter: F) -> Self
320 where
321 F: Fn(&LogItem) -> bool + Send + 'static,
322 {
323 self.filter = Some(Box::new(filter));
324 self
325 }
326
327 pub fn transformer<F>(mut self, transformer: F) -> Self
334 where
335 F: Fn(LogItem) -> LogItem + Send + 'static,
336 {
337 self.transformer = Some(Box::new(transformer));
338 self
339 }
340
341 pub fn start(self) -> (Logger, LoggerGuard) {
348 let (tx, mut rx) = unbounded_channel();
349 let (stop_tx, mut stop_rx) = unbounded_channel();
350 tracing_subscriber::fmt()
351 .with_max_level(self.level)
352 .with_writer(LogSender::new(tx.clone()))
353 .without_time()
354 .with_file(self.filename)
355 .with_line_number(self.line_number)
356 .json()
357 .init();
358
359 let join = thread::spawn(move || {
360 let handler = |v: Map<String, Value>| async {
361 if let Some(x) = &self.handler {
362 if !x(&v).await {
363 return;
364 }
365 }
366 let mut item = LogItem::from_json(v);
367 let time = item.fields.remove("_time").unwrap_or_default().as_i64();
368 if self.json {
369 item.time = time.unwrap_or_else(|| Utc::now().timestamp_micros()).into();
370 } else {
371 item.time = time
372 .map_or_else(Local::now, |v| {
373 DateTime::from_timestamp_micros(v)
374 .unwrap_or_default()
375 .into()
376 })
377 .format("%F %T%.6f")
378 .to_string()
379 .into();
380 }
381
382 if let Some(filter) = &self.filter {
383 if !filter(&item) {
384 return;
385 }
386 }
387 if let Some(transformer) = &self.transformer {
388 item = transformer(item);
389 }
390 let writer: Box<dyn io::Write> = if item.level <= Level::WARN {
391 Box::new(stderr())
392 } else {
393 Box::new(stdout())
394 };
395 if self.json {
396 let _ = (self.json_writer)(item, writer);
397 } else {
398 let _ = (self.color_writer)(item, writer);
399 }
400 };
401 block_on(async move {
402 loop {
403 select! {
404 Some(v) = rx.recv() => {
405 handler(v).await;
406 },
407 _ = stop_rx.recv() => {
408 while let Ok(v) = rx.try_recv(){
409 handler(v).await;
410 }
411 break;
412 }
413 }
414 }
415 })
416 });
417 (
418 Logger { tx },
419 LoggerGuard {
420 stop_tx,
421 join: Some(join),
422 },
423 )
424 }
425}