Skip to content
//! # 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);
    }
}