1use crate::{
12 counters::{
13 PROCESSED_STRUCT_LOG_COUNT, SENT_STRUCT_LOG_BYTES,
14 SENT_STRUCT_LOG_COUNT, STRUCT_LOG_PARSE_ERROR_COUNT,
15 STRUCT_LOG_QUEUE_ERROR_COUNT, STRUCT_LOG_SEND_ERROR_COUNT,
16 },
17 logger::Logger,
18 struct_log::TcpWriter,
19 Event, Filter, Level, LevelFilter, Metadata,
20};
21use backtrace::Backtrace;
22use chrono::{SecondsFormat, Utc};
23use once_cell::sync::Lazy;
24use parking_lot::{Mutex, RwLock};
25use pipe_logger_lib::{PipeLogger, PipeLoggerBuilder, RotateMethod};
26use serde::Serialize;
27use std::{
28 collections::BTreeMap,
29 env, fmt,
30 io::Write,
31 sync::{
32 mpsc::{self, Receiver, SyncSender},
33 Arc,
34 },
35 thread,
36};
37
38const RUST_LOG: &str = "RUST_LOG";
39pub const CHANNEL_SIZE: usize = 10000;
42const NUM_SEND_RETRIES: u8 = 1;
43
44#[derive(Debug, Serialize)]
46pub struct LogEntry {
47 #[serde(flatten)]
48 metadata: Metadata,
49 #[serde(skip_serializing_if = "Option::is_none")]
50 thread_name: Option<String>,
51 #[serde(skip_serializing_if = "Option::is_none")]
54 backtrace: Option<String>,
55 #[serde(skip_serializing_if = "Option::is_none")]
56 hostname: Option<&'static str>,
57 timestamp: String,
58 #[serde(skip_serializing_if = "BTreeMap::is_empty")]
59 data: BTreeMap<&'static str, serde_json::Value>,
60 #[serde(skip_serializing_if = "Option::is_none")]
61 message: Option<String>,
62}
63
64impl LogEntry {
65 fn new(event: &Event, thread_name: Option<&str>) -> Self {
66 use crate::{Key, Value, Visitor};
67
68 struct JsonVisitor<'a>(
69 &'a mut BTreeMap<&'static str, serde_json::Value>,
70 );
71
72 impl<'a> Visitor for JsonVisitor<'a> {
73 fn visit_pair(&mut self, key: Key, value: Value<'_>) {
74 let v = match value {
75 Value::Debug(d) => {
76 serde_json::Value::String(format!("{:?}", d))
77 }
78 Value::Display(d) => {
79 serde_json::Value::String(d.to_string())
80 }
81 Value::Serde(s) => match serde_json::to_value(s) {
82 Ok(value) => value,
83 Err(e) => {
84 eprintln!(
85 "error serializing structured log: {}",
86 e
87 );
88 return;
89 }
90 },
91 };
92
93 self.0.insert(key.as_str(), v);
94 }
95 }
96
97 let metadata = *event.metadata();
98 let thread_name = thread_name.map(ToOwned::to_owned);
99 let message = event.message().map(fmt::format);
100
101 static HOSTNAME: Lazy<Option<String>> = Lazy::new(|| {
102 hostname::get()
103 .ok()
104 .and_then(|name| name.into_string().ok())
105 });
106
107 let hostname = HOSTNAME.as_deref();
108
109 let backtrace = match metadata.level() {
110 Level::Error => {
111 let mut backtrace = Backtrace::new();
112 let mut frames = backtrace.frames().to_vec();
113 if frames.len() > 3 {
114 frames.drain(0..3); }
118 backtrace = frames.into();
119 Some(format!("{:?}", backtrace))
120 }
121 _ => None,
122 };
123
124 let mut data = BTreeMap::new();
125 for schema in event.keys_and_values() {
126 schema.visit(&mut JsonVisitor(&mut data));
127 }
128
129 Self {
130 metadata,
131 thread_name,
132 backtrace,
133 hostname,
134 timestamp: Utc::now().to_rfc3339_opts(SecondsFormat::Micros, true),
135 data,
136 message,
137 }
138 }
139
140 pub fn metadata(&self) -> &Metadata { &self.metadata }
141
142 pub fn thread_name(&self) -> Option<&str> { self.thread_name.as_deref() }
143
144 pub fn backtrace(&self) -> Option<&str> { self.backtrace.as_deref() }
145
146 pub fn hostname(&self) -> Option<&str> { self.hostname.as_deref() }
147
148 pub fn timestamp(&self) -> &str { self.timestamp.as_str() }
149
150 pub fn data(&self) -> &BTreeMap<&'static str, serde_json::Value> {
151 &self.data
152 }
153
154 pub fn message(&self) -> Option<&str> { self.message.as_deref() }
155}
156
157pub struct DiemLoggerBuilder {
159 channel_size: usize,
160 level: Level,
161 remote_level: Level,
162 address: Option<String>,
163 printer: Option<Box<dyn Writer>>,
164 is_async: bool,
165 custom_format: Option<fn(&LogEntry) -> Result<String, fmt::Error>>,
166}
167
168impl DiemLoggerBuilder {
169 #[allow(clippy::new_without_default)]
170 pub fn new() -> Self {
171 Self {
172 channel_size: CHANNEL_SIZE,
173 level: Level::Info,
174 remote_level: Level::Debug,
175 address: None,
176 printer: Some(Box::new(StderrWriter)),
177 is_async: false,
178 custom_format: None,
179 }
180 }
181
182 pub fn address(&mut self, address: String) -> &mut Self {
183 self.address = Some(address);
184 self
185 }
186
187 pub fn read_env(&mut self) -> &mut Self {
188 if let Ok(address) = env::var("STRUCT_LOG_TCP_ADDR") {
189 self.address(address);
190 }
191 self
192 }
193
194 pub fn level(&mut self, level: Level) -> &mut Self {
195 self.level = level;
196 self
197 }
198
199 pub fn remote_level(&mut self, level: Level) -> &mut Self {
200 self.remote_level = level;
201 self
202 }
203
204 pub fn channel_size(&mut self, channel_size: usize) -> &mut Self {
205 self.channel_size = channel_size;
206 self
207 }
208
209 pub fn printer(
210 &mut self, printer: Box<dyn Writer + Send + Sync + 'static>,
211 ) -> &mut Self {
212 self.printer = Some(printer);
213 self
214 }
215
216 pub fn is_async(&mut self, is_async: bool) -> &mut Self {
217 self.is_async = is_async;
218 self
219 }
220
221 pub fn custom_format(
222 &mut self, format: fn(&LogEntry) -> Result<String, fmt::Error>,
223 ) -> &mut Self {
224 self.custom_format = Some(format);
225 self
226 }
227
228 pub fn init(&mut self) { self.build(); }
229
230 pub fn build(&mut self) -> Arc<DiemLogger> {
231 let filter = {
232 let local_filter = {
233 let mut filter_builder = Filter::builder();
234
235 if env::var(RUST_LOG).is_ok() {
236 filter_builder.with_env(RUST_LOG);
237 } else {
238 filter_builder.filter_level(self.level.into());
239 }
240
241 filter_builder.build()
242 };
243 let remote_filter = {
244 let mut filter_builder = Filter::builder();
245
246 if self.is_async && self.address.is_some() {
247 filter_builder.filter_level(self.remote_level.into());
248 } else {
249 filter_builder.filter_level(LevelFilter::Off);
250 }
251
252 filter_builder.build()
253 };
254
255 DiemFilter {
256 local_filter,
257 remote_filter,
258 }
259 };
260
261 let logger = if self.is_async {
262 let (sender, receiver) = mpsc::sync_channel(self.channel_size);
263 let logger = Arc::new(DiemLogger {
264 sender: Some(sender),
265 printer: None,
266 filter: RwLock::new(filter),
267 formatter: self.custom_format.take().unwrap_or(default_format),
268 });
269 let service = LoggerService {
270 receiver,
271 address: self.address.clone(),
272 printer: self.printer.take(),
273 facade: logger.clone(),
274 };
275
276 thread::spawn(move || service.run());
277 logger
278 } else {
279 Arc::new(DiemLogger {
280 sender: None,
281 printer: self.printer.take(),
282 filter: RwLock::new(filter),
283 formatter: self.custom_format.take().unwrap_or(default_format),
284 })
285 };
286
287 crate::logger::set_global_logger(logger.clone());
288 logger
289 }
290}
291
292struct DiemFilter {
294 local_filter: Filter,
296 remote_filter: Filter,
298}
299
300impl DiemFilter {
301 fn enabled(&self, metadata: &Metadata) -> bool {
302 self.local_filter.enabled(metadata)
303 || self.remote_filter.enabled(metadata)
304 }
305}
306
307pub struct DiemLogger {
308 sender: Option<SyncSender<LoggerServiceEvent>>,
309 printer: Option<Box<dyn Writer>>,
310 filter: RwLock<DiemFilter>,
311 pub(crate) formatter: fn(&LogEntry) -> Result<String, fmt::Error>,
312}
313
314impl DiemLogger {
315 pub fn builder() -> DiemLoggerBuilder { DiemLoggerBuilder::new() }
316
317 #[allow(clippy::new_ret_no_self)]
318 pub fn new() -> DiemLoggerBuilder { Self::builder() }
319
320 pub fn init_for_testing() {
321 if env::var(RUST_LOG).is_err() {
322 return;
323 }
324
325 Self::builder()
326 .is_async(false)
327 .printer(Box::new(StderrWriter))
328 .build();
329 }
330
331 pub fn set_filter(&self, filter: Filter) {
332 self.filter.write().local_filter = filter;
333 }
334
335 pub fn set_remote_filter(&self, filter: Filter) {
336 self.filter.write().remote_filter = filter;
337 }
338
339 fn send_entry(&self, entry: LogEntry) {
340 if let Some(printer) = &self.printer {
341 let s = (self.formatter)(&entry).expect("Unable to format");
342 printer.write(s);
343 }
344
345 if let Some(sender) = &self.sender {
346 if let Err(e) = sender.try_send(LoggerServiceEvent::LogEntry(entry))
347 {
348 STRUCT_LOG_QUEUE_ERROR_COUNT.inc();
349 eprintln!("Failed to send structured log: {}", e);
350 }
351 }
352 }
353}
354
355impl Logger for DiemLogger {
356 fn enabled(&self, metadata: &Metadata) -> bool {
357 self.filter.read().enabled(metadata)
358 }
359
360 fn record(&self, event: &Event) {
361 let entry = LogEntry::new(event, ::std::thread::current().name());
362
363 self.send_entry(entry)
364 }
365
366 fn flush(&self) {
367 if let Some(sender) = &self.sender {
368 let (oneshot_sender, oneshot_receiver) = mpsc::sync_channel(1);
369 sender
370 .send(LoggerServiceEvent::Flush(oneshot_sender))
371 .unwrap();
372 oneshot_receiver.recv().unwrap();
373 }
374 }
375}
376
377enum LoggerServiceEvent {
378 LogEntry(LogEntry),
379 Flush(SyncSender<()>),
380}
381
382struct LoggerService {
385 receiver: Receiver<LoggerServiceEvent>,
386 address: Option<String>,
387 printer: Option<Box<dyn Writer>>,
388 facade: Arc<DiemLogger>,
389}
390
391impl LoggerService {
392 pub fn run(mut self) {
393 let mut writer = self.address.take().map(TcpWriter::new);
394
395 for event in self.receiver {
396 match event {
397 LoggerServiceEvent::LogEntry(entry) => {
398 PROCESSED_STRUCT_LOG_COUNT.inc();
399
400 if let Some(printer) = &self.printer {
401 if self
402 .facade
403 .filter
404 .read()
405 .local_filter
406 .enabled(&entry.metadata)
407 {
408 let s = (self.facade.formatter)(&entry)
409 .expect("Unable to format");
410 printer.write(s)
411 }
412 }
413
414 if let Some(writer) = &mut writer {
415 if self
416 .facade
417 .filter
418 .read()
419 .remote_filter
420 .enabled(&entry.metadata)
421 {
422 Self::write_to_logstash(writer, entry);
423 }
424 }
425 }
426 LoggerServiceEvent::Flush(sender) => {
427 let _ = sender.send(());
431 }
432 }
433 }
434 }
435
436 fn write_to_logstash(stream: &mut TcpWriter, mut entry: LogEntry) {
439 if entry.message.is_none() {
442 entry.message = Some(serde_json::to_string(&entry.data).unwrap());
443 }
444
445 let message = if let Ok(json) = serde_json::to_string(&entry) {
446 json
447 } else {
448 STRUCT_LOG_PARSE_ERROR_COUNT.inc();
449 return;
450 };
451
452 let message = message + "\n";
453 let bytes = message.as_bytes();
454 let message_length = bytes.len();
455
456 let mut result = stream.write_all(bytes);
460 for _ in 0..NUM_SEND_RETRIES {
461 if result.is_ok() {
462 break;
463 } else {
464 result = stream.write_all(bytes);
465 }
466 }
467
468 if let Err(e) = result {
469 STRUCT_LOG_SEND_ERROR_COUNT.inc();
470 eprintln!(
471 "[Logging] Error while sending data to logstash({}): {}",
472 stream.endpoint(),
473 e
474 );
475 } else {
476 SENT_STRUCT_LOG_COUNT.inc();
477 SENT_STRUCT_LOG_BYTES.inc_by(message_length as u64);
478 }
479 }
480}
481
482pub trait Writer: Send + Sync {
484 fn write(&self, log: String);
486}
487
488struct StderrWriter;
490
491impl Writer for StderrWriter {
492 fn write(&self, log: String) {
494 eprintln!("{}", log);
495 }
496}
497
498pub struct FileWriter {
500 log_file: RwLock<std::fs::File>,
501}
502
503impl FileWriter {
504 pub fn new(log_file: std::path::PathBuf) -> Self {
505 let file = std::fs::OpenOptions::new()
506 .append(true)
507 .create(true)
508 .open(log_file)
509 .expect("Unable to open log file");
510 Self {
511 log_file: RwLock::new(file),
512 }
513 }
514}
515
516impl Writer for FileWriter {
517 fn write(&self, log: String) {
519 if let Err(err) = writeln!(self.log_file.write(), "{}", log) {
520 eprintln!("Unable to write to log file: {}", err.to_string());
521 }
522 }
523}
524
525pub struct RollingFileWriter {
526 log_file: Mutex<PipeLogger>,
527}
528
529impl RollingFileWriter {
530 pub fn new(
531 log_path: std::path::PathBuf, count: usize, size_mb: usize,
532 ) -> Self {
533 let mut builder = PipeLoggerBuilder::new(log_path);
534 builder
535 .set_compress(true)
536 .set_rotate(Some(RotateMethod::FileSize(
537 size_mb as u64 * 1_000_000,
538 )))
539 .set_count(Some(count));
540 let log_file = builder.build().unwrap();
541 Self {
542 log_file: Mutex::new(log_file),
543 }
544 }
545}
546
547impl Writer for RollingFileWriter {
548 fn write(&self, log: String) {
549 if let Err(e) = self.log_file.lock().write_line(&log) {
550 eprintln!("Unable to write to log file: {}", e.to_string());
551 }
552 }
553}
554
555fn default_format(entry: &LogEntry) -> Result<String, fmt::Error> {
561 use std::fmt::Write;
562
563 let mut w = String::new();
564 write!(w, "{}", entry.timestamp)?;
565
566 if let Some(thread_name) = &entry.thread_name {
567 write!(w, " [{}]", thread_name)?;
568 }
569
570 write!(
571 w,
572 " {} {}",
573 entry.metadata.level(),
574 entry.metadata.location()
575 )?;
576
577 if let Some(message) = &entry.message {
578 write!(w, " {}", message)?;
579 }
580
581 if !entry.data.is_empty() {
582 write!(w, " {}", serde_json::to_string(&entry.data).unwrap())?;
583 }
584
585 Ok(w)
586}
587
588#[cfg(test)]
589mod tests {
590 use super::LogEntry;
591 use crate::{
592 debug, error, info, logger::Logger, trace, warn, Event, Key, KeyValue,
593 Level, Metadata, Schema, Value, Visitor,
594 };
595 use chrono::{DateTime, Utc};
596 use serde_json::Value as JsonValue;
597 use std::{
598 sync::{
599 mpsc::{self, Receiver, SyncSender},
600 Arc,
601 },
602 thread,
603 };
604
605 #[derive(serde::Serialize)]
606 #[serde(rename_all = "snake_case")]
607 enum Enum {
608 FooBar,
609 }
610
611 struct TestSchema<'a> {
612 foo: usize,
613 bar: &'a Enum,
614 }
615
616 impl Schema for TestSchema<'_> {
617 fn visit(&self, visitor: &mut dyn Visitor) {
618 visitor.visit_pair(Key::new("foo"), Value::from_serde(&self.foo));
619 visitor.visit_pair(Key::new("bar"), Value::from_serde(&self.bar));
620 }
621 }
622
623 struct LogStream(SyncSender<LogEntry>);
624
625 impl LogStream {
626 fn new() -> (Self, Receiver<LogEntry>) {
627 let (sender, receiver) = mpsc::sync_channel(1024);
628 (Self(sender), receiver)
629 }
630 }
631
632 impl Logger for LogStream {
633 fn enabled(&self, metadata: &Metadata) -> bool {
634 metadata.level() <= Level::Debug
635 }
636
637 fn record(&self, event: &Event) {
638 let entry = LogEntry::new(event, ::std::thread::current().name());
639 self.0.send(entry).unwrap();
640 }
641
642 fn flush(&self) {}
643 }
644
645 fn set_test_logger() -> Receiver<LogEntry> {
646 let (logger, receiver) = LogStream::new();
647 let logger = Arc::new(logger);
648 crate::logger::set_global_logger(logger);
649 receiver
650 }
651
652 #[test]
655 fn basic() {
656 let receiver = set_test_logger();
657 let number = 12345;
658
659 let before = Utc::now();
661 info!(
662 TestSchema {
663 foo: 5,
664 bar: &Enum::FooBar
665 },
666 test = true,
667 category = "name",
668 KeyValue::new("display", Value::from_display(&number)),
669 "This is a log"
670 );
671 let after = Utc::now();
672
673 let entry = receiver.recv().unwrap();
674
675 assert_eq!(entry.metadata.level(), Level::Info);
677 assert_eq!(
678 entry.metadata.target(),
679 module_path!().split("::").next().unwrap()
680 );
681 assert_eq!(entry.metadata.module_path(), module_path!());
682 assert_eq!(entry.metadata.file(), file!());
683 assert_eq!(entry.message.as_deref(), Some("This is a log"));
684 assert!(entry.backtrace.is_none());
685
686 let timestamp = DateTime::parse_from_rfc3339(&entry.timestamp).unwrap();
688 let timestamp: DateTime<Utc> = DateTime::from(timestamp);
689 assert!(before <= timestamp && timestamp <= after);
690
691 assert_eq!(entry.data.get("foo").and_then(JsonValue::as_u64), Some(5));
693 assert_eq!(
694 entry.data.get("bar").and_then(JsonValue::as_str),
695 Some("foo_bar")
696 );
697 assert_eq!(
698 entry.data.get("display").and_then(JsonValue::as_str),
699 Some(format!("{}", number)).as_deref(),
700 );
701 assert_eq!(
702 entry.data.get("test").and_then(JsonValue::as_bool),
703 Some(true),
704 );
705 assert_eq!(
706 entry.data.get("category").and_then(JsonValue::as_str),
707 Some("name"),
708 );
709
710 error!("This is an error log");
712 let entry = receiver.recv().unwrap();
713 assert!(entry.backtrace.is_some());
714
715 trace!("trace");
719 debug!("debug");
720 info!("info");
721 warn!("warn");
722 error!("error");
723
724 let levels = &[Level::Debug, Level::Info, Level::Warn, Level::Error];
725
726 for level in levels {
727 let entry = receiver.recv().unwrap();
728 assert_eq!(entry.metadata.level(), *level);
729 }
730
731 let handler = thread::Builder::new()
733 .name("named thread".into())
734 .spawn(|| info!("thread"))
735 .unwrap();
736
737 handler.join().unwrap();
738 let entry = receiver.recv().unwrap();
739 assert_eq!(entry.thread_name.as_deref(), Some("named thread"));
740
741 let debug_struct = DebugStruct {};
743 let display_struct = DisplayStruct {};
744
745 error!(identifier = ?debug_struct, "Debug test");
746 error!(identifier = ?debug_struct, other = "value", "Debug2 test");
747 error!(identifier = %display_struct, "Display test");
748 error!(identifier = %display_struct, other = "value", "Display2 test");
749 error!("Literal" = ?debug_struct, "Debug test");
750 error!("Literal" = ?debug_struct, other = "value", "Debug test");
751 error!("Literal" = %display_struct, "Display test");
752 error!("Literal" = %display_struct, other = "value", "Display2 test");
753 error!("Literal" = %display_struct, other = "value", identifier = ?debug_struct, "Mixed test");
754 }
755
756 struct DebugStruct {}
757
758 impl std::fmt::Debug for DebugStruct {
759 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
760 write!(f, "DebugStruct!")
761 }
762 }
763
764 struct DisplayStruct {}
765
766 impl std::fmt::Display for DisplayStruct {
767 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
768 write!(f, "DisplayStruct!")
769 }
770 }
771}