Skip to main content

actix_cloud/
logger.rs

1//! Provide logger feature.
2//! The inner library uses [tracing](https://crates.io/crates/tracing).
3//! See their documents for how to log in the program.
4//!
5//! This wrapper makes it thread safe, even for FFI libraries.
6//! You can use it everywhere and freely.
7use 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        // We do not buffer output.
105        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    /// Get logger sender.
122    pub fn sender(&self) -> UnboundedSender<Map<String, Value>> {
123        self.tx.clone()
124    }
125
126    /// Init tracing logger.
127    /// A new subscriber will be registered.
128    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
145/// Keep this guard alive when you use the logger.
146/// No more logs will be record after dropping the guard.
147///
148/// # Warning
149/// When dropping the guard, it will wait for the logger thread to exit.
150pub 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    /// Return colored string of `level`.
184    ///
185    /// - TRACE/DEBUG => Magenta
186    /// - INFO => Green
187    /// - WARN => Yellow
188    /// - ERROR => Red
189    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    /// Create new logger instance.
240    /// Default is colorful writer, INFO level, no filename and line number.
241    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    /// Use custom json writer.
256    ///
257    /// # Warning
258    /// Do not perform heavy workloads, it can block other logs!
259    pub fn json_writer(mut self, writer: WriterFn) -> Self {
260        self.json_writer = writer;
261        self
262    }
263
264    /// Use custom colorful writer.
265    ///
266    /// # Warning
267    /// Do not perform heavy workloads, it can block other logs!
268    pub fn color_writer(mut self, writer: WriterFn) -> Self {
269        self.color_writer = writer;
270        self
271    }
272
273    /// Use json format writer.
274    pub fn json(mut self) -> Self {
275        self.json = true;
276        self
277    }
278
279    /// Set log level.
280    pub fn level(mut self, level: Level) -> Self {
281        self.level = level;
282        self
283    }
284
285    /// Enable filename in the log.
286    pub fn filename(mut self) -> Self {
287        self.filename = true;
288        self
289    }
290
291    /// Enable line number in the log.
292    pub fn line_number(mut self) -> Self {
293        self.line_number = true;
294        self
295    }
296
297    /// Customize the handler.
298    ///
299    /// The customized handler will be invoked first, even before the filter.
300    /// When the return value is false, further handler will be skipped.
301    /// Otherwise, normal log hander will still be invoked.
302    ///
303    /// # Warning
304    /// Do not perform heavy workloads, it can block other logs!
305    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    /// Customize the filter. Filter out unwanted logs.
314    ///
315    /// When the filter function return false, no logs will be sent to the transformer.
316    ///
317    /// # Warning
318    /// Do not perform heavy workloads, it can block other logs!
319    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    /// Customize the transformer. Change the logs on the fly.
328    ///
329    /// After this function, LogItem will be sent to the corresponding writer.
330    ///
331    /// # Warning
332    /// Do not perform heavy workloads, it can block other logs!
333    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    /// Start logger.
342    /// This method will spawn a new thread to print the log.
343    ///
344    /// You should call this method only once for the entire program.
345    /// For FFI library, you need to call this method once in the library code and keep the return values alive.
346    /// Then customize the [Self::handler] and send output back to the main program.
347    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}