//! # void-logger
//!
//! A logging backend for the standard `log` crate facade that formats and prints
//! log lines to the console (stdout/stderr) while simultaneously queuing and sending
//! them in compressed batches to a Void ingestion server.
//!
//! ## Example Setup
//!
//! ```rust
//! use std::time::Duration;
//! use void_logger::{VoidLoggerBuilder, ConsoleOutput};
//!
//! fn main() {
//! VoidLoggerBuilder::new("http://127.0.0.1:9180/ingest")
//! .service_name("my-service")
//! .max_level(log::LevelFilter::Info)
//! .console_output(ConsoleOutput::Split)
//! .init()
//! .expect("Failed to initialize VoidLogger");
//!
//! log::info!("Hello, world!");
//! }
//! ```
use chrono::Utc;
use log::{Log, Metadata, Record};
use serde_json::Value;
use std::collections::HashMap;
use std::sync::mpsc::{Receiver, RecvTimeoutError, SyncSender, sync_channel};
use std::time::{Duration, Instant};
use void_common::LogLine;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ContentEncoding {
Identity,
Gzip,
Zstd,
Deflate,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ContentType {
Json,
MsgPack,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ConsoleOutput {
None,
Stdout,
Stderr,
Split,
}
#[allow(dead_code)]
enum LogCommand {
Log(LogLine),
Flush(std::sync::mpsc::Sender<()>),
Shutdown,
}
pub struct VoidLogger {
tx: SyncSender<LogCommand>,
max_level: log::LevelFilter,
service_name: String,
static_metadata: HashMap<String, Value>,
console_output: ConsoleOutput,
}
impl VoidLogger {
fn new(builder: VoidLoggerBuilder) -> Self {
let (tx, rx) = sync_channel(builder.buffer_size);
let url = builder.url;
let content_encoding = builder.content_encoding;
let content_type = builder.content_type;
let batch_size = builder.batch_size;
let batch_timeout = builder.batch_timeout;
let api_key = builder.api_key;
std::thread::spawn(move || {
worker_loop(
rx,
url,
content_encoding,
content_type,
batch_size,
batch_timeout,
api_key,
);
});
Self {
tx,
max_level: builder.max_level,
service_name: builder.service_name,
static_metadata: builder.static_metadata,
console_output: builder.console_output,
}
}
fn record_to_log_line(&self, record: &Record) -> LogLine {
let mut metadata = self.static_metadata.clone();
if let Some(file) = record.file() {
metadata.insert("file".to_string(), Value::String(file.to_string()));
}
if let Some(line) = record.line() {
metadata.insert("line".to_string(), Value::Number(line.into()));
}
metadata.insert(
"target".to_string(),
Value::String(record.target().to_string()),
);
if let Some(module) = record.module_path() {
metadata.insert("module_path".to_string(), Value::String(module.to_string()));
}
LogLine {
timestamp: Utc::now(),
severity: record.level().to_string(),
message: record.args().to_string(),
service: self.service_name.clone(),
metadata,
}
}
}
impl Log for VoidLogger {
fn enabled(&self, metadata: &Metadata) -> bool {
metadata.level() <= self.max_level
}
fn log(&self, record: &Record) {
if self.enabled(record.metadata()) {
let timestamp = Utc::now().format("%Y-%m-%dT%H:%M:%S%.3fZ");
match self.console_output {
ConsoleOutput::None => {}
ConsoleOutput::Stdout => {
println!(
"[{}] {:<5} [{}] {}",
timestamp,
record.level(),
record.target(),
record.args()
);
}
ConsoleOutput::Stderr => {
eprintln!(
"[{}] {:<5} [{}] {}",
timestamp,
record.level(),
record.target(),
record.args()
);
}
ConsoleOutput::Split => {
let formatted = format!(
"[{}] {:<5} [{}] {}",
timestamp,
record.level(),
record.target(),
record.args()
);
if record.level() <= log::Level::Warn {
eprintln!("{}", formatted);
} else {
println!("{}", formatted);
}
}
}
let log_line = self.record_to_log_line(record);
// Non-blocking send: if buffer is full, print a warning to stderr and discard log
if let Err(_e) = self.tx.try_send(LogCommand::Log(log_line)) {
// To avoid spamming stderr, we could ignore, but print warning in debug or optionally
#[cfg(debug_assertions)]
eprintln!(
"VoidLogger warning: log buffer full or worker thread died, log discarded"
);
}
}
}
fn flush(&self) {
let (tx, rx) = std::sync::mpsc::channel();
if self.tx.send(LogCommand::Flush(tx)).is_ok() {
let _ = rx.recv_timeout(Duration::from_secs(5));
}
}
}
pub struct VoidLoggerBuilder {
url: String,
service_name: String,
max_level: log::LevelFilter,
batch_size: usize,
batch_timeout: Duration,
buffer_size: usize,
static_metadata: HashMap<String, Value>,
content_encoding: ContentEncoding,
content_type: ContentType,
console_output: ConsoleOutput,
api_key: Option<String>,
}
impl VoidLoggerBuilder {
pub fn new(url: impl Into<String>) -> Self {
Self {
url: url.into(),
service_name: "unknown".to_string(),
max_level: log::LevelFilter::Info,
batch_size: 100,
batch_timeout: Duration::from_secs(1),
buffer_size: 10000,
static_metadata: HashMap::new(),
content_encoding: ContentEncoding::Identity,
content_type: ContentType::Json,
console_output: ConsoleOutput::Split,
api_key: None,
}
}
pub fn service_name(mut self, service_name: impl Into<String>) -> Self {
self.service_name = service_name.into();
self
}
pub fn max_level(mut self, max_level: log::LevelFilter) -> Self {
self.max_level = max_level;
self
}
pub fn batch_size(mut self, batch_size: usize) -> Self {
self.batch_size = batch_size;
self
}
pub fn batch_timeout(mut self, batch_timeout: Duration) -> Self {
self.batch_timeout = batch_timeout;
self
}
pub fn buffer_size(mut self, buffer_size: usize) -> Self {
self.buffer_size = buffer_size;
self
}
pub fn static_metadata(mut self, metadata: HashMap<String, Value>) -> Self {
self.static_metadata = metadata;
self
}
pub fn add_metadata(mut self, key: impl Into<String>, value: impl Into<Value>) -> Self {
self.static_metadata.insert(key.into(), value.into());
self
}
pub fn content_encoding(mut self, encoding: ContentEncoding) -> Self {
self.content_encoding = encoding;
self
}
pub fn content_type(mut self, content_type: ContentType) -> Self {
self.content_type = content_type;
self
}
pub fn console_output(mut self, console_output: ConsoleOutput) -> Self {
self.console_output = console_output;
self
}
pub fn api_key(mut self, api_key: impl Into<String>) -> Self {
self.api_key = Some(api_key.into());
self
}
pub fn init(self) -> Result<(), log::SetLoggerError> {
let max_level = self.max_level;
let logger = VoidLogger::new(self);
log::set_boxed_logger(Box::new(logger))?;
log::set_max_level(max_level);
Ok(())
}
}
fn worker_loop(
rx: Receiver<LogCommand>,
url: String,
content_encoding: ContentEncoding,
content_type: ContentType,
batch_size: usize,
batch_timeout: Duration,
api_key: Option<String>,
) {
let mut buffer = Vec::with_capacity(batch_size);
let mut first_log_time: Option<Instant> = None;
loop {
let timeout = match first_log_time {
None => Duration::from_secs(3600 * 24), // effectively infinity when buffer is empty
Some(start) => {
let elapsed = start.elapsed();
if elapsed >= batch_timeout {
Duration::from_secs(0)
} else {
batch_timeout - elapsed
}
}
};
let msg = if timeout.is_zero() {
Err(RecvTimeoutError::Timeout)
} else {
rx.recv_timeout(timeout)
};
match msg {
Ok(LogCommand::Log(log)) => {
if buffer.is_empty() {
first_log_time = Some(Instant::now());
}
buffer.push(log);
if buffer.len() >= batch_size {
send_batch(
&url,
&buffer,
content_encoding,
content_type,
api_key.as_deref(),
);
buffer.clear();
first_log_time = None;
}
}
Ok(LogCommand::Flush(ack_tx)) => {
if !buffer.is_empty() {
send_batch(
&url,
&buffer,
content_encoding,
content_type,
api_key.as_deref(),
);
buffer.clear();
first_log_time = None;
}
let _ = ack_tx.send(());
}
Ok(LogCommand::Shutdown) => {
if !buffer.is_empty() {
send_batch(
&url,
&buffer,
content_encoding,
content_type,
api_key.as_deref(),
);
}
break;
}
Err(RecvTimeoutError::Timeout) => {
if !buffer.is_empty() {
send_batch(
&url,
&buffer,
content_encoding,
content_type,
api_key.as_deref(),
);
buffer.clear();
first_log_time = None;
}
}
Err(RecvTimeoutError::Disconnected) => {
if !buffer.is_empty() {
send_batch(
&url,
&buffer,
content_encoding,
content_type,
api_key.as_deref(),
);
}
break;
}
}
}
}
fn send_batch(
url: &str,
batch: &[LogLine],
encoding: ContentEncoding,
content_type: ContentType,
api_key: Option<&str>,
) {
if batch.is_empty() {
return;
}
// 1. Serialize
let serialized = match content_type {
ContentType::Json => match serde_json::to_vec(batch) {
Ok(bytes) => bytes,
Err(e) => {
eprintln!("VoidLogger error: failed to serialize batch to JSON: {}", e);
return;
}
},
#[cfg(feature = "msgpack")]
ContentType::MsgPack => match rmp_serde::to_vec(batch) {
Ok(bytes) => bytes,
Err(e) => {
eprintln!(
"VoidLogger error: failed to serialize batch to MsgPack: {}",
e
);
return;
}
},
#[cfg(not(feature = "msgpack"))]
ContentType::MsgPack => {
eprintln!("VoidLogger error: MsgPack feature is not enabled. Falling back to JSON.");
match serde_json::to_vec(batch) {
Ok(bytes) => bytes,
Err(e) => {
eprintln!("VoidLogger error: failed to serialize batch to JSON: {}", e);
return;
}
}
}
};
// 2. Compress
let (compressed, content_encoding_header) = match encoding {
ContentEncoding::Identity => (serialized, "identity"),
#[cfg(feature = "gzip")]
ContentEncoding::Gzip => {
use flate2::Compression;
use flate2::write::GzEncoder;
use std::io::Write;
let mut encoder = GzEncoder::new(Vec::new(), Compression::default());
if encoder.write_all(&serialized).is_ok() {
if let Ok(bytes) = encoder.finish() {
(bytes, "gzip")
} else {
eprintln!("VoidLogger error: failed to gzip compress logs");
(serialized, "identity")
}
} else {
eprintln!("VoidLogger error: failed to write to gzip encoder");
(serialized, "identity")
}
}
#[cfg(not(feature = "gzip"))]
ContentEncoding::Gzip => {
eprintln!("VoidLogger error: gzip feature is not enabled. Sending uncompressed.");
(serialized, "identity")
}
#[cfg(feature = "zstd")]
ContentEncoding::Zstd => match zstd::encode_all(&serialized[..], 0) {
Ok(bytes) => (bytes, "zstd"),
Err(e) => {
eprintln!("VoidLogger error: failed to zstd compress logs: {}", e);
(serialized, "identity")
}
},
#[cfg(not(feature = "zstd"))]
ContentEncoding::Zstd => {
eprintln!("VoidLogger error: zstd feature is not enabled. Sending uncompressed.");
(serialized, "identity")
}
ContentEncoding::Deflate => {
#[cfg(feature = "gzip")]
// We can use flate2 for deflate since gzip feature includes flate2
{
use flate2::Compression;
use flate2::write::DeflateEncoder;
use std::io::Write;
let mut encoder = DeflateEncoder::new(Vec::new(), Compression::default());
if encoder.write_all(&serialized).is_ok() {
if let Ok(bytes) = encoder.finish() {
(bytes, "deflate")
} else {
eprintln!("VoidLogger error: failed to deflate compress logs");
(serialized, "identity")
}
} else {
eprintln!("VoidLogger error: failed to write to deflate encoder");
(serialized, "identity")
}
}
#[cfg(not(feature = "gzip"))]
{
eprintln!(
"VoidLogger error: gzip/flate2 feature is not enabled. Sending uncompressed."
);
(serialized, "identity")
}
}
};
// 3. Send
let content_type_header = match content_type {
ContentType::Json => "application/json",
ContentType::MsgPack => "application/msgpack",
};
let mut request = minreq::post(url)
.with_header("Content-Type", content_type_header)
.with_header("Content-Encoding", content_encoding_header)
.with_header("X-Content-Encoding", content_encoding_header)
.with_timeout(5);
if let Some(key) = api_key {
request = request.with_header("Authorization", format!("Bearer {}", key));
}
let response = request.with_body(compressed).send();
match response {
Ok(res) => {
if res.status_code < 200 || res.status_code >= 300 {
eprintln!(
"VoidLogger error: Server returned status {} for ingest: {}",
res.status_code,
res.as_str().unwrap_or_default()
);
}
}
Err(e) => {
eprintln!("VoidLogger error: failed to send logs to server: {}", e);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::io::Read;
use std::net::TcpListener;
use std::sync::mpsc::channel;
#[test]
fn test_logger_integration() {
// 1. Start a simple TCP mock server to capture the POST request
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let url = format!("http://127.0.0.1:{}/ingest", port);
let (tx, rx) = channel();
// Run the listener in a background thread to accept one request
std::thread::spawn(move || {
if let Ok((mut stream, _)) = listener.accept() {
let mut buffer = vec![0; 16384];
if let Ok(bytes_read) = stream.read(&mut buffer) {
let req_data = buffer[..bytes_read].to_vec();
let _ = tx.send(req_data);
// Respond with HTTP 200 OK
let response = b"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: 21\r\n\r\n{\"status\":\"success\"}";
use std::io::Write;
let _ = stream.write_all(response);
}
}
});
// 2. Build and initialize our logger
let builder = VoidLoggerBuilder::new(url)
.service_name("test-service")
.max_level(log::LevelFilter::Info)
.batch_size(1) // send immediately
.add_metadata("env", "testing")
.api_key("super-secret-logger-key");
let logger = VoidLogger::new(builder);
// 3. Log a line
let record = log::Record::builder()
.args(format_args!("Hello Void Logger!"))
.level(log::Level::Info)
.target("test_target")
.file(Some("lib.rs"))
.line(Some(42))
.build();
logger.log(&record);
logger.flush();
// 4. Receive and parse the captured request
let raw_req = rx
.recv_timeout(Duration::from_secs(5))
.expect("Mock server did not receive request");
let req_str = String::from_utf8_lossy(&raw_req);
// Assert HTTP headers
assert!(req_str.contains("POST /ingest HTTP/1.1"));
assert!(req_str.contains("Authorization: Bearer super-secret-logger-key"));
let req_str_lower = req_str.to_lowercase();
assert!(req_str_lower.contains("content-type: application/json"));
// Extract body (after \r\n\r\n)
let parts: Vec<&str> = req_str.split("\r\n\r\n").collect();
assert!(parts.len() >= 2);
let body = parts[1];
// Parse JSON body to Vec<LogLine>
let logs: Vec<LogLine> =
serde_json::from_str(body).expect("Failed to parse logs JSON from request body");
assert_eq!(logs.len(), 1);
let log = &logs[0];
assert_eq!(log.severity, "INFO");
assert_eq!(log.message, "Hello Void Logger!");
assert_eq!(log.service, "test-service");
assert_eq!(
log.metadata.get("env").unwrap().as_str().unwrap(),
"testing"
);
assert_eq!(
log.metadata.get("file").unwrap().as_str().unwrap(),
"lib.rs"
);
assert_eq!(log.metadata.get("line").unwrap().as_u64().unwrap(), 42);
assert_eq!(
log.metadata.get("target").unwrap().as_str().unwrap(),
"test_target"
);
}
#[test]
fn test_logger_gzip_compression() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let url = format!("http://127.0.0.1:{}/ingest", port);
let (tx, rx) = channel();
std::thread::spawn(move || {
if let Ok((mut stream, _)) = listener.accept() {
let mut buffer = vec![0; 16384];
if let Ok(bytes_read) = stream.read(&mut buffer) {
let req_data = buffer[..bytes_read].to_vec();
let _ = tx.send(req_data);
let response = b"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: 21\r\n\r\n{\"status\":\"success\"}";
use std::io::Write;
let _ = stream.write_all(response);
}
}
});
let builder = VoidLoggerBuilder::new(url)
.service_name("compress-service")
.max_level(log::LevelFilter::Info)
.batch_size(1)
.content_encoding(ContentEncoding::Gzip);
let logger = VoidLogger::new(builder);
let record = log::Record::builder()
.args(format_args!("Compressed Log"))
.level(log::Level::Info)
.build();
logger.log(&record);
logger.flush();
let raw_req = rx
.recv_timeout(Duration::from_secs(5))
.expect("Mock server did not receive request");
let req_str = String::from_utf8_lossy(&raw_req);
let req_str_lower = req_str.to_lowercase();
assert!(
req_str_lower.contains("content-encoding: gzip")
|| req_str_lower.contains("x-content-encoding: gzip")
);
}
#[test]
fn test_console_output_configurations() {
let builder_stdout =
VoidLoggerBuilder::new("http://127.0.0.1:0").console_output(ConsoleOutput::Stdout);
let logger_stdout = VoidLogger::new(builder_stdout);
let record = log::Record::builder()
.args(format_args!("Test stdout print"))
.level(log::Level::Info)
.build();
logger_stdout.log(&record);
let builder_none =
VoidLoggerBuilder::new("http://127.0.0.1:0").console_output(ConsoleOutput::None);
let logger_none = VoidLogger::new(builder_none);
logger_none.log(&record);
let builder_stderr =
VoidLoggerBuilder::new("http://127.0.0.1:0").console_output(ConsoleOutput::Stderr);
let logger_stderr = VoidLogger::new(builder_stderr);
logger_stderr.log(&record);
let builder_split =
VoidLoggerBuilder::new("http://127.0.0.1:0").console_output(ConsoleOutput::Split);
let logger_split = VoidLogger::new(builder_split);
let info_record = log::Record::builder()
.args(format_args!("Test split stdout print"))
.level(log::Level::Info)
.build();
logger_split.log(&info_record);
let warn_record = log::Record::builder()
.args(format_args!("Test split stderr print"))
.level(log::Level::Warn)
.build();
logger_split.log(&warn_record);
}
}