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/// One structured log entry, handed to the writers (and `filter`/`transformer`).
30#[serde_as]
31#[derive(Serialize, Deserialize, Debug)]
32pub struct LogItem {
33    /// Timestamp: microseconds since epoch (JSON output) or formatted string
34    /// (color output). Overridable through the reserved `_time` field.
35    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    /// Extra key-value pairs attached to the event (`log.*` fields are internal).
42    #[serde(skip_serializing_if = "Map::is_empty")]
43    pub fields: Map<String, Value>,
44    /// Key-value pairs of the current span (`http.*`/`otel.*`/`name` are internal).
45    #[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        // We do not buffer output.
110        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
120/// Build and register the tracing subscriber that forwards JSON events into `tx`.
121fn 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/// Handle to the running logger; clone and pass around (e.g. store in `GlobalState`)
138/// to re-init the subscriber or feed raw log maps from plugins.
139#[derive(Clone)]
140pub struct Logger {
141    tx: UnboundedSender<Map<String, Value>>,
142}
143
144impl Logger {
145    /// Get logger sender.
146    pub fn sender(&self) -> UnboundedSender<Map<String, Value>> {
147        self.tx.clone()
148    }
149
150    /// Init tracing logger.
151    /// A new subscriber will be registered.
152    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
162/// Writer callback: render one [`LogItem`] to the given sink (stdout/stderr).
163pub type WriterFn = Box<dyn Fn(LogItem, Box<dyn Write>) -> Result<()> + Send>;
164/// Filter callback: return `false` to drop a log entry before the transformer.
165pub type FilterFn = Box<dyn Fn(&LogItem) -> bool + Send>;
166/// Transformer callback: rewrite a log entry before writing.
167pub type TransformerFn = Box<dyn Fn(LogItem) -> LogItem + Send>;
168/// Handler callback: invoked first on the raw JSON map; return `false` to stop
169/// further processing of the entry.
170pub type HandlerFn = Box<dyn Fn(&Map<String, Value>) -> Pin<Box<dyn Future<Output = bool>>> + Send>;
171
172/// Keep this guard alive when you use the logger.
173/// No more logs will be record after dropping the guard.
174///
175/// # Warning
176/// When dropping the guard, it will wait for the logger thread to exit.
177pub 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    /// Return colored string of `level`.
211    ///
212    /// - TRACE/DEBUG => Magenta
213    /// - INFO => Green
214    /// - WARN => Yellow
215    /// - ERROR => Red
216    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    /// Create new logger instance.
265    /// Default is colorful writer, INFO level, no filename and line number.
266    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    /// Use custom json writer.
281    ///
282    /// # Warning
283    /// Do not perform heavy workloads, it can block other logs!
284    pub fn json_writer(mut self, writer: WriterFn) -> Self {
285        self.json_writer = writer;
286        self
287    }
288
289    /// Use custom colorful writer.
290    ///
291    /// # Warning
292    /// Do not perform heavy workloads, it can block other logs!
293    pub fn color_writer(mut self, writer: WriterFn) -> Self {
294        self.color_writer = writer;
295        self
296    }
297
298    /// Use json format writer.
299    pub fn json(mut self) -> Self {
300        self.json = true;
301        self
302    }
303
304    /// Set log level.
305    pub fn level(mut self, level: Level) -> Self {
306        self.level = level;
307        self
308    }
309
310    /// Enable filename in the log.
311    pub fn filename(mut self) -> Self {
312        self.filename = true;
313        self
314    }
315
316    /// Enable line number in the log.
317    pub fn line_number(mut self) -> Self {
318        self.line_number = true;
319        self
320    }
321
322    /// Customize the handler.
323    ///
324    /// The customized handler will be invoked first, even before the filter.
325    /// When the return value is false, further handler will be skipped.
326    /// Otherwise, normal log handler will still be invoked.
327    ///
328    /// # Warning
329    /// Do not perform heavy workloads, it can block other logs!
330    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    /// Customize the filter. Filter out unwanted logs.
339    ///
340    /// When the filter function return false, no logs will be sent to the transformer.
341    ///
342    /// # Warning
343    /// Do not perform heavy workloads, it can block other logs!
344    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    /// Customize the transformer. Change the logs on the fly.
353    ///
354    /// After this function, LogItem will be sent to the corresponding writer.
355    ///
356    /// # Warning
357    /// Do not perform heavy workloads, it can block other logs!
358    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    /// Start logger.
367    /// This method will spawn a new thread to print the log.
368    ///
369    /// You should call this method only once for the entire program.
370    /// For FFI library, you need to call this method once in the library code and keep the return values alive.
371    /// Then customize the [Self::handler] and send output back to the main program.
372    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}