use chrono::{DateTime, Utc};
use serde::Serialize;
use std::fmt::{Debug, Display};
use std::sync::Arc;
mod json;
mod valid;
pub use futures::future::BoxFuture;
pub use valid::*;
pub struct Logger {
max_level: Level,
sink: Arc<dyn Sink + Send + Sync>,
}
impl Logger {
pub fn new(max_level: Level, sink: Arc<dyn Sink + Send + Sync>) -> Self {
Self {
max_level: max_level,
sink: sink,
}
}
pub fn log<T: Clone + Serialize>(
&self,
mode: SinkMode,
level: Level,
entry: Entry<T>,
) -> SinkAcknowledgment {
self.log_expensive(mode, level, || entry)
}
pub fn log_expensive<F, T: Clone + Serialize>(
&self,
mode: SinkMode,
level: Level,
func: F,
) -> SinkAcknowledgment
where
F: FnOnce() -> Entry<T>,
{
if level < self.max_level {
SinkAcknowledgment::NotPerformed
} else {
let entry = func();
match mode {
SinkMode::Blocking => {
SinkAcknowledgment::Completed(self.sink.sink_blocking(entry.json(level)))
}
SinkMode::Awaitable => {
SinkAcknowledgment::Awaitable(self.sink.sink(entry.json(level)))
}
}
}
}
pub fn log_blocking<T: Clone + Serialize>(
&self,
level: Level,
entry: Entry<T>,
) -> Option<std::io::Result<()>> {
match self.log(SinkMode::Blocking, level, entry) {
SinkAcknowledgment::Completed(result) => Some(result),
SinkAcknowledgment::NotPerformed => None,
SinkAcknowledgment::Awaitable(_) => Some(Err(new_mode_ack_mismatch_err())),
}
}
pub fn log_async<T: Clone + Serialize>(
&self,
level: Level,
entry: Entry<T>,
) -> Option<BoxFuture<std::io::Result<()>>> {
match self.log(SinkMode::Awaitable, level, entry) {
SinkAcknowledgment::NotPerformed => None,
SinkAcknowledgment::Completed(_) => Some(Box::pin(futures::future::ready(Err(
new_mode_ack_mismatch_err(),
)))),
SinkAcknowledgment::Awaitable(fut) => Some(fut),
}
}
pub fn log_expensive_blocking<F, T: Clone + Serialize>(
&self,
level: Level,
func: F,
) -> Option<std::io::Result<()>>
where
F: FnOnce() -> Entry<T>,
{
match self.log_expensive(SinkMode::Blocking, level, func) {
SinkAcknowledgment::NotPerformed => None,
SinkAcknowledgment::Completed(result) => Some(result),
SinkAcknowledgment::Awaitable(_) => Some(Err(new_mode_ack_mismatch_err())),
}
}
pub fn log_expensive_async<F, T: Clone + Serialize>(
&self,
level: Level,
func: F,
) -> Option<BoxFuture<std::io::Result<()>>>
where
F: FnOnce() -> Entry<T>,
{
match self.log_expensive(SinkMode::Awaitable, level, func) {
SinkAcknowledgment::NotPerformed => None,
SinkAcknowledgment::Completed(_) => Some(Box::pin(futures::future::ready(Err(
new_mode_ack_mismatch_err(),
)))),
SinkAcknowledgment::Awaitable(fut) => Some(fut),
}
}
pub fn max_level(&self) -> Level {
self.max_level
}
pub fn set_max_level(&mut self, level: Level) {
self.max_level = level;
}
}
const MODE_ACK_MISMATCH_ERR_MESSAGE: &'static str =
"sink acknowledgment does not match mode requested";
fn new_mode_ack_mismatch_err() -> std::io::Error {
std::io::Error::new(std::io::ErrorKind::Other, MODE_ACK_MISMATCH_ERR_MESSAGE)
}
#[derive(Copy, Clone, Debug, Eq, PartialEq, Ord, PartialOrd, Hash)]
pub enum SinkMode {
Blocking,
Awaitable,
}
pub enum SinkAcknowledgment<'a> {
NotPerformed,
Completed(std::io::Result<()>),
Awaitable(BoxFuture<'a, std::io::Result<()>>),
}
#[derive(Copy, Clone, Debug, Eq, PartialEq, Ord, PartialOrd, Hash, Serialize)]
pub enum Level {
TRACE,
DEBUG,
INFO,
WARN,
ERROR,
FATAL,
}
impl Display for Level {
fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
Debug::fmt(self, f)
}
}
#[derive(Clone, Serialize)]
pub struct Entry<T: Clone + Serialize> {
#[serde(rename = "msg")]
message: T,
#[serde(rename = "ts")]
#[serde(skip_serializing_if = "Option::is_none")]
created: Option<DateTime<Utc>>,
#[serde(skip_serializing_if = "Option::is_none")]
#[serde(serialize_with = "json::as_base64")]
#[serde(rename = "spid")]
span_id: Option<Vec<u8>>,
#[serde(skip_serializing_if = "Option::is_none")]
tags: Option<Tags>,
}
impl<T: Clone + Serialize> Entry<T> {
pub fn new(message: T, timestamp: bool, span_id: Option<Vec<u8>>, tags: Option<Tags>) -> Self {
Self {
message: message,
created: if timestamp { Some(Utc::now()) } else { None },
span_id: span_id,
tags: tags,
}
}
pub fn json(self, level: Level) -> String {
let mut log_line = json::serialize::<T>(self, level);
log_line.push('\n');
log_line
}
}
pub trait Sink {
fn sink_blocking(&self, entry: String) -> std::io::Result<()>;
fn sink(&self, entry: String) -> BoxFuture<std::io::Result<()>>;
}
#[cfg(test)]
mod tests {
use super::*;
use futures::executor::block_on;
#[test]
fn stringify_enum() {
assert_eq!("WARN", format!("{}", Level::WARN));
assert_eq!("TRACE", format!("{}", Level::TRACE));
}
struct DummySink;
impl Sink for DummySink {
fn sink_blocking(&self, _entry: String) -> std::io::Result<()> {
Ok(())
}
fn sink(&self, _entry: String) -> BoxFuture<std::io::Result<()>> {
Box::pin(async { Ok(()) })
}
}
#[test]
fn log_expensive() {
let dummy_sink = DummySink;
let logger = Logger {
max_level: Level::INFO,
sink: Arc::new(dummy_sink),
};
let mut i = 400;
let entry = Entry::new("asffdf", false, None, None);
logger.log_expensive(SinkMode::Blocking, Level::TRACE, || {
i += 1;
entry.clone()
});
assert_eq!(i, 400);
logger.log_expensive(SinkMode::Blocking, Level::FATAL, || {
i += 1;
entry.clone()
});
assert_eq!(i, 401);
logger.log_expensive(SinkMode::Blocking, Level::INFO, || {
i += 1;
entry.clone()
});
assert_eq!(i, 402);
}
#[test]
fn log_shortcuts() {
let dummy_sink = DummySink;
let logger = Logger {
max_level: Level::INFO,
sink: Arc::new(dummy_sink),
};
let entry = Entry::new("asffdf", false, None, None);
assert!(logger
.log_blocking(Level::ERROR, entry.clone())
.unwrap()
.is_ok());
assert!(block_on(logger.log_async(Level::ERROR, entry.clone()).unwrap()).is_ok());
let mut i = 400;
assert!(logger
.log_expensive_blocking(Level::TRACE, || {
i += 1;
entry.clone()
})
.is_none());
assert_eq!(i, 400);
assert!(block_on(
logger
.log_expensive_async(Level::ERROR, || {
i += 1;
entry.clone()
})
.unwrap(),
)
.is_ok());
assert_eq!(i, 401);
}
}