use super::LogSink;
use crate::ConsoleSinkConfig;
use crate::DataMasker;
use crate::InklogError;
use crate::LogRecord;
use crate::LogTemplate;
use crate::support::processing::OutputFormat;
use async_trait::async_trait;
use is_terminal::IsTerminal;
use owo_colors::OwoColorize;
use std::fmt;
use std::io::{self, Write};
use std::sync::{Arc, Mutex};
pub struct ConsoleSink {
config: ConsoleSinkConfig,
writer: Arc<Mutex<Box<dyn Write + Send>>>,
template: LogTemplate,
masker: DataMasker,
}
impl fmt::Debug for ConsoleSink {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("ConsoleSink")
.field("config", &self.config)
.field("template", &self.template)
.finish()
}
}
impl ConsoleSink {
pub fn new(config: ConsoleSinkConfig, template: LogTemplate) -> Self {
Self {
config,
writer: Arc::new(Mutex::new(Box::new(io::stdout()))),
template,
masker: DataMasker::new(),
}
}
fn write_record<W: Write>(
&self,
writer: &mut W,
record: &LogRecord,
use_color: bool,
) -> io::Result<()> {
if self.config.output_format == OutputFormat::Json {
let json = serde_json::to_string(record).map_err(io::Error::other)?;
writeln!(writer, "{}", json)
} else {
let formatted_message = self.template.render(record);
if use_color {
writeln!(
writer,
"{}",
self.apply_color(&formatted_message, &record.level)
)
} else {
writeln!(writer, "{}", formatted_message)
}
}
}
fn apply_color(&self, message: &str, level: &str) -> String {
match level {
"ERROR" | "error" => message.red().to_string(),
"WARN" | "warn" => message.yellow().to_string(),
"INFO" | "info" => message.green().to_string(),
"DEBUG" | "debug" => message.blue().to_string(),
"TRACE" | "trace" => message.magenta().to_string(),
_ => message.green().to_string(),
}
}
fn should_colorize(&self, is_stderr: bool) -> bool {
if self.config.output_format == OutputFormat::Json {
return false;
}
if !self.config.colored {
return false;
}
if std::env::var("NO_COLOR").is_ok() {
return false;
}
if let Ok(val) = std::env::var("CLICOLOR_FORCE")
&& val != "0"
{
return true;
}
if let Ok(term) = std::env::var("TERM")
&& term == "dumb"
{
return false;
}
if is_stderr {
io::stderr().is_terminal()
} else {
io::stdout().is_terminal()
}
}
}
#[async_trait]
impl LogSink for ConsoleSink {
async fn write(&self, record: &LogRecord) -> Result<(), InklogError> {
let masked_record = if self.config.masking_enabled {
let mut masked = record.clone();
masked.message = self.masker.mask(&record.message);
self.masker.mask_hashmap(&mut masked.fields);
masked
} else {
record.clone()
};
let is_stderr = self
.config
.stderr_levels
.contains(&masked_record.level.to_lowercase());
let use_color = self.should_colorize(is_stderr);
if is_stderr {
let mut stderr = io::stderr();
self.write_record(&mut stderr, &masked_record, use_color)
.map_err(InklogError::IoError)?;
} else {
let mut writer = self
.writer
.lock()
.map_err(|_| InklogError::IoError(io::Error::other("Lock poisoned")))?;
self.write_record(&mut *writer, &masked_record, use_color)
.map_err(InklogError::IoError)?;
}
Ok(())
}
async fn flush(&self) -> Result<(), InklogError> {
let mut writer = self
.writer
.lock()
.map_err(|_| InklogError::IoError(io::Error::other("Lock poisoned")))?;
writer.flush().map_err(InklogError::IoError)?;
io::stderr().flush().map_err(InklogError::IoError)
}
fn is_healthy(&self) -> bool {
true
}
async fn shutdown(&self) -> Result<(), InklogError> {
self.flush().await
}
}
impl Clone for ConsoleSink {
fn clone(&self) -> Self {
Self {
config: self.config.clone(),
writer: Arc::clone(&self.writer),
template: self.template.clone(),
masker: DataMasker::new(),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ConsoleSinkConfig;
use serial_test::serial;
use std::env;
fn get_sink() -> ConsoleSink {
ConsoleSink::new(
ConsoleSinkConfig {
enabled: true,
colored: true,
..Default::default()
},
LogTemplate::default(),
)
}
#[test]
#[serial]
fn test_no_color_env() {
let sink = get_sink();
unsafe {
env::set_var("NO_COLOR", "1");
}
assert!(!sink.should_colorize(false));
unsafe {
env::remove_var("NO_COLOR");
}
}
#[test]
#[serial]
fn test_force_color_env() {
let sink = get_sink();
unsafe {
env::remove_var("NO_COLOR");
}
unsafe {
env::set_var("CLICOLOR_FORCE", "1");
}
assert!(sink.should_colorize(false));
unsafe {
env::remove_var("CLICOLOR_FORCE");
}
}
#[test]
#[serial]
fn test_term_dumb() {
let sink = get_sink();
unsafe {
env::remove_var("NO_COLOR");
}
unsafe {
env::set_var("TERM", "dumb");
}
unsafe {
env::remove_var("CLICOLOR_FORCE");
}
assert!(!sink.should_colorize(false));
unsafe {
env::remove_var("TERM");
}
}
#[test]
#[serial]
fn test_config_disabled() {
let mut sink = get_sink();
unsafe {
env::remove_var("NO_COLOR");
}
sink.config.colored = false;
unsafe {
env::set_var("CLICOLOR_FORCE", "1");
} assert!(!sink.should_colorize(false));
unsafe {
env::remove_var("CLICOLOR_FORCE");
}
}
#[test]
fn test_console_sink_new() {
let config = ConsoleSinkConfig {
enabled: true,
colored: true,
..Default::default()
};
let template = LogTemplate::default();
let sink = ConsoleSink::new(config, template);
assert!(sink.config.enabled);
}
#[test]
fn test_console_sink_disabled() {
let config = ConsoleSinkConfig {
enabled: false,
colored: true,
..Default::default()
};
let template = LogTemplate::default();
let sink = ConsoleSink::new(config, template);
assert!(!sink.config.enabled);
}
#[test]
#[serial]
fn test_should_colorize_defaults() {
unsafe {
env::remove_var("CLICOLOR_FORCE");
}
unsafe {
env::remove_var("TERM");
}
unsafe {
env::set_var("NO_COLOR", "1");
}
let config = ConsoleSinkConfig {
enabled: true,
colored: true,
..Default::default()
};
let template = LogTemplate::default();
let sink = ConsoleSink::new(config, template);
let result = sink.should_colorize(false);
assert!(
!result,
"should_colorize should return false when NO_COLOR is set"
);
unsafe {
env::remove_var("NO_COLOR");
}
}
#[test]
fn test_should_colorize_when_allowed() {
let sink = get_sink();
let colored = sink.apply_color("test message", "ERROR");
assert!(colored.contains("test message"));
}
#[test]
fn test_apply_color_info() {
let sink = get_sink();
let colored = sink.apply_color("test message", "INFO");
assert!(colored.contains("test message"));
}
#[test]
fn test_apply_color_unknown() {
let sink = get_sink();
let colored = sink.apply_color("test message", "UNKNOWN");
assert!(colored.contains("test message"));
}
#[derive(Default, Clone)]
struct TestWriter {
buf: Arc<Mutex<Vec<u8>>>,
}
impl Write for TestWriter {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
self.buf.lock().unwrap().write(buf)
}
fn flush(&mut self) -> io::Result<()> {
self.buf.lock().unwrap().flush()
}
}
impl TestWriter {
fn output(&self) -> String {
let buf = self.buf.lock().unwrap();
String::from_utf8_lossy(&buf).to_string()
}
fn is_empty(&self) -> bool {
self.buf.lock().unwrap().is_empty()
}
}
struct FailingWriter;
impl Write for FailingWriter {
fn write(&mut self, _buf: &[u8]) -> io::Result<usize> {
Err(io::Error::other("write failed"))
}
fn flush(&mut self) -> io::Result<()> {
Err(io::Error::other("flush failed"))
}
}
fn make_record(level: &str, message: &str) -> LogRecord {
LogRecord {
level: level.to_string(),
message: message.to_string(),
target: "test::module".to_string(),
..Default::default()
}
}
fn sink_with_test_writer(config: ConsoleSinkConfig) -> (ConsoleSink, TestWriter) {
let writer = TestWriter::default();
let mut sink = ConsoleSink::new(config, LogTemplate::default());
sink.writer = Arc::new(Mutex::new(Box::new(writer.clone())));
(sink, writer)
}
#[test]
fn test_apply_color_warn() {
let sink = get_sink();
let colored = sink.apply_color("warn message", "WARN");
assert!(colored.contains("warn message"));
}
#[test]
fn test_apply_color_debug() {
let sink = get_sink();
let colored = sink.apply_color("debug message", "DEBUG");
assert!(colored.contains("debug message"));
}
#[test]
fn test_apply_color_trace() {
let sink = get_sink();
let colored = sink.apply_color("trace message", "TRACE");
assert!(colored.contains("trace message"));
}
#[test]
fn test_apply_color_lowercase_levels() {
let sink = get_sink();
for level in &["error", "warn", "info", "debug", "trace"] {
let colored = sink.apply_color("payload", level);
assert!(colored.contains("payload"), "level {} lost message", level);
}
}
#[test]
#[serial]
fn test_apply_color_emits_ansi_codes() {
owo_colors::set_override(true);
let sink = get_sink();
let red = sink.apply_color("msg", "ERROR");
let yellow = sink.apply_color("msg", "WARN");
let green = sink.apply_color("msg", "INFO");
let blue = sink.apply_color("msg", "DEBUG");
let magenta = sink.apply_color("msg", "TRACE");
owo_colors::unset_override();
assert!(
red.contains("\x1b[31m"),
"ERROR must be red, got: {:?}",
red
);
assert!(
yellow.contains("\x1b[33m"),
"WARN must be yellow, got: {:?}",
yellow
);
assert!(
green.contains("\x1b[32m"),
"INFO must be green, got: {:?}",
green
);
assert!(
blue.contains("\x1b[34m"),
"DEBUG must be blue, got: {:?}",
blue
);
assert!(
magenta.contains("\x1b[35m"),
"TRACE must be magenta, got: {:?}",
magenta
);
}
#[test]
fn test_write_record_without_color() {
let sink = get_sink();
let mut buf: Vec<u8> = Vec::new();
let record = make_record("INFO", "hello world");
sink.write_record(&mut buf, &record, false).unwrap();
let output = String::from_utf8(buf).unwrap();
assert!(output.contains("[INFO]"));
assert!(output.contains("test::module"));
assert!(output.contains("hello world"));
assert!(output.ends_with('\n'), "writeln should append newline");
assert!(!output.contains('\x1b'));
}
#[test]
fn test_write_record_with_color_error() {
let sink = get_sink();
let mut buf: Vec<u8> = Vec::new();
let record = make_record("ERROR", "boom");
sink.write_record(&mut buf, &record, true).unwrap();
let output = String::from_utf8(buf).unwrap();
assert!(output.contains("boom"));
assert!(output.contains("[ERROR]"));
assert!(output.ends_with('\n'));
}
#[test]
fn test_write_record_with_color_warn() {
let sink = get_sink();
let mut buf: Vec<u8> = Vec::new();
let record = make_record("WARN", "careful");
sink.write_record(&mut buf, &record, true).unwrap();
let output = String::from_utf8(buf).unwrap();
assert!(output.contains("careful"));
assert!(output.contains("[WARN]"));
}
#[test]
fn test_write_record_with_color_debug() {
let sink = get_sink();
let mut buf: Vec<u8> = Vec::new();
let record = make_record("DEBUG", "details");
sink.write_record(&mut buf, &record, true).unwrap();
let output = String::from_utf8(buf).unwrap();
assert!(output.contains("details"));
assert!(output.contains("[DEBUG]"));
}
#[test]
fn test_write_record_with_color_trace() {
let sink = get_sink();
let mut buf: Vec<u8> = Vec::new();
let record = make_record("TRACE", "verbose");
sink.write_record(&mut buf, &record, true).unwrap();
let output = String::from_utf8(buf).unwrap();
assert!(output.contains("verbose"));
assert!(output.contains("[TRACE]"));
}
#[test]
fn test_write_record_with_color_unknown_level() {
let sink = get_sink();
let mut buf: Vec<u8> = Vec::new();
let record = make_record("FATAL", "critical");
sink.write_record(&mut buf, &record, true).unwrap();
let output = String::from_utf8(buf).unwrap();
assert!(output.contains("critical"));
assert!(output.contains("[FATAL]"));
assert!(output.ends_with('\n'));
}
#[test]
fn test_write_record_lowercase_level_with_color() {
let sink = get_sink();
let mut buf: Vec<u8> = Vec::new();
let record = make_record("error", "lowercase boom");
sink.write_record(&mut buf, &record, true).unwrap();
let output = String::from_utf8(buf).unwrap();
assert!(output.contains("lowercase boom"));
assert!(output.contains("[error]"));
}
#[test]
fn test_write_record_propagates_write_error() {
let sink = get_sink();
let mut writer = FailingWriter;
let record = make_record("INFO", "will fail");
let result = sink.write_record(&mut writer, &record, false);
assert!(result.is_err());
let err = result.unwrap_err();
assert!(
err.to_string().contains("write failed"),
"error should carry the underlying message, got: {}",
err
);
}
#[test]
#[serial]
fn test_should_colorize_clicolor_force_zero_falls_through() {
unsafe {
env::remove_var("NO_COLOR");
}
unsafe {
env::set_var("CLICOLOR_FORCE", "0");
}
unsafe {
env::set_var("TERM", "dumb");
}
let sink = get_sink();
assert!(!sink.should_colorize(false));
unsafe {
env::remove_var("CLICOLOR_FORCE");
}
unsafe {
env::remove_var("TERM");
}
}
#[test]
#[serial]
fn test_should_colorize_stderr_path_with_force() {
unsafe {
env::remove_var("NO_COLOR");
}
unsafe {
env::set_var("CLICOLOR_FORCE", "1");
}
let sink = get_sink();
assert!(sink.should_colorize(true));
unsafe {
env::remove_var("CLICOLOR_FORCE");
}
}
#[test]
#[serial]
fn test_should_colorize_term_not_dumb_falls_through() {
unsafe {
env::remove_var("NO_COLOR");
}
unsafe {
env::set_var("TERM", "xterm-256color");
}
unsafe {
env::set_var("CLICOLOR_FORCE", "1");
}
let sink = get_sink();
assert!(sink.should_colorize(false));
unsafe {
env::remove_var("TERM");
}
unsafe {
env::remove_var("CLICOLOR_FORCE");
}
}
#[tokio::test]
async fn test_log_sink_write_stdout_no_masking() {
let config = ConsoleSinkConfig {
enabled: true,
colored: false,
masking_enabled: false,
..Default::default()
};
let (sink, writer) = sink_with_test_writer(config);
let record = make_record("INFO", "plain message");
sink.write(&record).await.unwrap();
let output = writer.output();
assert!(output.contains("plain message"));
assert!(output.contains("[INFO]"));
assert!(output.ends_with('\n'));
}
#[tokio::test]
async fn test_log_sink_write_with_masking_redacts_sensitive_data() {
let config = ConsoleSinkConfig {
enabled: true,
colored: false,
masking_enabled: true,
..Default::default()
};
let (sink, writer) = sink_with_test_writer(config);
let record = make_record("INFO", "email=test@example.com");
sink.write(&record).await.unwrap();
let output = writer.output();
assert!(
!output.contains("test@example.com"),
"masked output must not contain the original email, got: {}",
output
);
assert!(output.contains('@'), "masked email should retain @");
}
#[tokio::test]
async fn test_log_sink_write_stderr_level_writes_to_stderr_not_stdout() {
let config = ConsoleSinkConfig {
enabled: true,
colored: false,
stderr_levels: vec!["error".to_string()],
..Default::default()
};
let (sink, writer) = sink_with_test_writer(config);
let record = make_record("ERROR", "stderr-only message");
sink.write(&record).await.unwrap();
assert!(
writer.is_empty(),
"stdout writer must be empty when level routes to stderr"
);
}
#[tokio::test]
async fn test_log_sink_write_warn_routes_to_stderr_by_default() {
let config = ConsoleSinkConfig {
enabled: true,
colored: false,
..Default::default()
};
let (sink, writer) = sink_with_test_writer(config);
let record = make_record("WARN", "warning via stderr");
sink.write(&record).await.unwrap();
assert!(
writer.is_empty(),
"WARN should route to stderr by default, not stdout"
);
}
#[tokio::test]
async fn test_log_sink_write_info_routes_to_stdout_by_default() {
let config = ConsoleSinkConfig {
enabled: true,
colored: false,
..Default::default()
};
let (sink, writer) = sink_with_test_writer(config);
let record = make_record("INFO", "info via stdout");
sink.write(&record).await.unwrap();
let output = writer.output();
assert!(output.contains("info via stdout"));
}
#[tokio::test]
async fn test_log_sink_write_case_insensitive_stderr_match() {
let config = ConsoleSinkConfig {
enabled: true,
colored: false,
stderr_levels: vec!["error".to_string()],
..Default::default()
};
let (sink, writer) = sink_with_test_writer(config);
let record = make_record("ERROR", "uppercase level");
sink.write(&record).await.unwrap();
assert!(
writer.is_empty(),
"uppercase ERROR must match lowercase stderr_levels"
);
}
#[tokio::test]
async fn test_log_sink_flush_succeeds() {
let (sink, _writer) = sink_with_test_writer(ConsoleSinkConfig::default());
assert!(sink.flush().await.is_ok());
}
#[test]
fn test_log_sink_is_healthy_always_true() {
let sink = get_sink();
assert!(sink.is_healthy());
}
#[tokio::test]
async fn test_log_sink_shutdown_flushes_without_error() {
let (sink, _writer) = sink_with_test_writer(ConsoleSinkConfig::default());
assert!(sink.shutdown().await.is_ok());
}
#[test]
fn test_console_sink_clone_preserves_config() {
let config = ConsoleSinkConfig {
enabled: true,
colored: true,
stderr_levels: vec!["error".to_string(), "warn".to_string()],
masking_enabled: true,
output_format: Default::default(),
};
let sink = ConsoleSink::new(config, LogTemplate::default());
let cloned = sink.clone();
assert!(cloned.config.enabled);
assert!(cloned.config.colored);
assert!(cloned.config.masking_enabled);
assert_eq!(cloned.config.stderr_levels, vec!["error", "warn"]);
assert!(sink.is_healthy());
assert!(cloned.is_healthy());
}
#[tokio::test]
async fn test_console_sink_clone_shares_writer_buffer() {
let config = ConsoleSinkConfig {
enabled: true,
colored: false,
..Default::default()
};
let (sink, writer) = sink_with_test_writer(config);
let cloned = sink.clone();
let record = make_record("INFO", "written via clone");
cloned.write(&record).await.unwrap();
let output = writer.output();
assert!(
output.contains("written via clone"),
"clone shares writer via Arc, output should be visible"
);
}
#[test]
fn test_console_sink_debug_format() {
let sink = get_sink();
let debug_str = format!("{:?}", sink);
assert!(debug_str.contains("ConsoleSink"));
assert!(debug_str.contains("config"));
assert!(debug_str.contains("template"));
assert!(!debug_str.contains("writer"));
assert!(!debug_str.contains("masker"));
}
#[test]
fn test_write_record_with_color_info() {
let sink = get_sink();
let mut buf: Vec<u8> = Vec::new();
let record = make_record("INFO", "info colored output");
sink.write_record(&mut buf, &record, true).unwrap();
let output = String::from_utf8(buf).unwrap();
assert!(output.contains("info colored output"));
assert!(output.contains("[INFO]"));
assert!(
output.ends_with('\n'),
"writeln should append newline, got: {:?}",
output
);
}
#[test]
fn test_write_record_with_color_propagates_write_error() {
let sink = get_sink();
let mut writer = FailingWriter;
let record = make_record("ERROR", "colored but fails");
let result = sink.write_record(&mut writer, &record, true);
assert!(
result.is_err(),
"write_record with color should propagate error"
);
let err = result.unwrap_err();
assert!(
err.to_string().contains("write failed"),
"error should carry the underlying message, got: {}",
err
);
}
#[test]
fn test_write_record_with_color_all_levels_propagate_error() {
let sink = get_sink();
for level in &["ERROR", "WARN", "INFO", "DEBUG", "TRACE", "FATAL"] {
let mut writer = FailingWriter;
let record = make_record(level, "payload");
let result = sink.write_record(&mut writer, &record, true);
assert!(
result.is_err(),
"write_record with color should fail for level: {}",
level
);
}
}
}