use super::builder::{LoggerBuilder, LoggerDependencies};
use super::recovery::SinkControlMessage;
use super::workers::WorkerParams;
#[allow(unused_imports)]
use crate::ConsoleSinkConfig;
use crate::InklogError;
use crate::LogRecord;
use crate::LogTemplate;
use crate::domain::core::LoggerSubscriber;
use crate::integrations::Cache;
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
use crate::integrations::Database;
use crate::support::io::ConsoleSink;
use crate::support::io::FileSink;
use crate::support::io::LogSink;
use crate::support::processing::RateLimiter;
use crate::validation::sanitize::LogSanitizer;
use crate::{FileSinkConfig, InklogConfig};
use crate::{HealthStatus, Metrics};
use crate::{LogAdapter, LogLogger};
use crossbeam_channel::{Sender, bounded};
#[allow(unused_imports)]
use std::path::Path;
use std::path::PathBuf;
use std::string::ToString;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use tracing::error;
#[cfg(feature = "http")]
use tracing::info;
use tracing_subscriber::prelude::*;
pub struct LoggerManager {
#[allow(dead_code)]
config: InklogConfig,
sender: Sender<Arc<LogRecord>>,
console_sender: Sender<Arc<LogRecord>>,
shutdown_txs: Vec<Sender<()>>,
#[allow(dead_code)]
console_sink: Arc<Mutex<dyn LogSink>>,
metrics: Arc<Metrics>,
worker_handles: Mutex<Vec<tokio::task::JoinHandle<()>>>,
control_tx: Sender<SinkControlMessage>,
effective_capacity: Arc<AtomicUsize>,
#[cfg(feature = "http")]
http_server_handle: Mutex<Option<tokio::task::JoinHandle<()>>>,
cache: Option<Arc<dyn Cache>>,
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
database: Option<Arc<dyn Database>>,
}
impl LoggerManager {
pub async fn new() -> Result<Self, InklogError> {
Self::with_dependencies(LoggerDependencies::default()).await
}
pub async fn with_dependencies(deps: LoggerDependencies) -> Result<Self, InklogError> {
Self::build_with_deps(deps).await
}
async fn build_with_deps(deps: LoggerDependencies) -> Result<Self, InklogError> {
let config = if let Some(ref config_provider) = deps.config {
let mut config = InklogConfig::default();
if let Some(level) = config_provider.get_string("global.level") {
config.global.level = level;
}
if let Some(format) = config_provider.get_string("global.format") {
config.global.format = format;
}
if let Some(masking) = config_provider.get_bool("global.masking_enabled") {
config.global.masking_enabled = masking;
}
if let Some(fallback) = config_provider.get_bool("global.auto_fallback") {
config.global.auto_fallback = fallback;
}
if config_provider
.get_bool("file_sink.enabled")
.unwrap_or(false)
{
let path = config_provider
.get_string("file_sink.path")
.map(PathBuf::from)
.unwrap_or_default();
let max_size = config_provider
.get_string("file_sink.max_size")
.unwrap_or_else(|| "100MB".to_string());
let compress = config_provider
.get_bool("file_sink.compress")
.unwrap_or(true);
config.file_sink = Some(FileSinkConfig {
enabled: true,
path,
max_size,
compress,
..Default::default()
});
}
if config_provider
.get_bool("http_server.enabled")
.unwrap_or(false)
{
let host = config_provider
.get_string("http_server.host")
.unwrap_or_else(|| "127.0.0.1".to_string());
let port = config_provider
.get_int("http_server.port")
.map(|p| p as u16)
.unwrap_or(9090);
config.http_server = Some(crate::HttpServerConfig {
enabled: true,
host,
port,
..Default::default()
});
}
if let Some(threads) = config_provider.get_int("performance.worker_threads") {
config.performance.worker_threads = threads as usize;
}
if let Some(capacity) = config_provider.get_int("performance.channel_capacity") {
config.performance.channel_capacity = capacity as usize;
}
config
} else {
InklogConfig::load_sync().unwrap_or_else(|_| InklogConfig::default())
};
let cache = deps.cache;
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
let database = deps.database;
let (mut manager, _subscriber, _filter) = Self::build_detached(
config,
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
database.clone(),
)
.await?;
manager.cache = cache;
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
{
manager.database = database;
}
Ok(manager)
}
pub async fn with_config(config: InklogConfig) -> Result<Self, InklogError> {
#[cfg(feature = "http")]
tracing::info!(
event = "security_logger_initialized",
sinks = ?config.sinks_enabled(),
masking_enabled = config.global.masking_enabled,
"Logger manager initialized"
);
let (manager, subscriber, filter) = Self::build_detached(
config.clone(),
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
None,
)
.await?;
let registry = tracing_subscriber::registry().with(subscriber).with(filter);
if let Err(ref e) = registry.try_init() {
tracing::debug!(error = %e, "global subscriber already set; skipping inklog registry");
}
let log_adapter = LogAdapter::new(
manager.console_sender.clone(),
manager.sender.clone(),
manager.metrics.clone(),
);
let max_level = config
.global
.level
.parse::<tracing::Level>()
.unwrap_or(tracing::Level::INFO);
let log_level = match max_level {
tracing::Level::TRACE => log::LevelFilter::Trace,
tracing::Level::DEBUG => log::LevelFilter::Debug,
tracing::Level::INFO => log::LevelFilter::Info,
tracing::Level::WARN => log::LevelFilter::Warn,
tracing::Level::ERROR => log::LevelFilter::Error,
};
let log_logger = LogLogger::new(log_adapter, log_level);
if let Err(e) = log_logger.install() {
tracing::debug!(error = %e, "log crate logger already set; skipping inklog LogLogger");
}
#[cfg(feature = "http")]
if let Some(ref http_cfg) = config.http_server
&& http_cfg.enabled
&& let Err(e) = Self::start_http_server(
manager.metrics.clone(),
manager.sender.clone(),
manager.effective_capacity.clone(),
&manager.http_server_handle,
http_cfg,
)
.await
{
match http_cfg.error_mode {
crate::HttpErrorMode::Warn => {
let mut args = fluent_bundle::FluentArgs::new();
args.set("err", e.to_string());
tracing::warn!(
"{}",
crate::i18n::tr_args("config-http_startup_failed", args)
);
}
crate::HttpErrorMode::Strict => {
return Err(e);
}
}
}
Ok(manager)
}
pub async fn build_detached(
config: InklogConfig,
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
database: Option<Arc<dyn Database>>,
) -> Result<
(
Self,
LoggerSubscriber,
tracing_subscriber::filter::EnvFilter,
),
InklogError,
> {
let metrics = Arc::new(Metrics::new());
let (sender, receiver) = bounded(config.performance.channel_capacity);
let (console_sender, console_receiver) = bounded(config.performance.channel_capacity);
let (control_tx, control_rx) = bounded(10); let effective_capacity = Arc::new(AtomicUsize::new(config.performance.channel_capacity));
let console_sink: Arc<Mutex<dyn LogSink>> = Arc::new(Mutex::new(ConsoleSink::new(
config.console_sink.clone().unwrap_or_default(),
LogTemplate::new(&config.global.format),
)));
let mut subscriber =
LoggerSubscriber::new(console_sender.clone(), sender.clone(), metrics.clone());
if config.global.masking_enabled {
subscriber = subscriber.with_sanitizer(Arc::new(LogSanitizer::new()));
}
if let Some(rate) = config.performance.rate_limit {
subscriber = subscriber.with_rate_limiter(Arc::new(RateLimiter::new(rate)));
}
let level = config
.global
.level
.parse::<tracing::Level>()
.unwrap_or(tracing::Level::INFO);
let level_str = match level {
tracing::Level::TRACE => "trace",
tracing::Level::DEBUG => "debug",
tracing::Level::INFO => "info",
tracing::Level::WARN => "warn",
tracing::Level::ERROR => "error",
};
let filter = match std::env::var("RUST_LOG") {
Ok(val) if !val.is_empty() => {
tracing_subscriber::filter::EnvFilter::new(format!("{},{}", level_str, val))
}
_ => tracing_subscriber::filter::EnvFilter::new(level_str),
};
let error_sink_config = FileSinkConfig {
enabled: true,
path: PathBuf::from("logs/error.log"),
..Default::default()
};
let error_sink: Arc<Mutex<Option<Box<dyn LogSink>>>> =
Arc::new(Mutex::new(match FileSink::new(error_sink_config) {
Ok(sink) => Some(Box::new(sink) as Box<dyn LogSink>),
Err(e) => {
tracing::warn!(error = %e, "Failed to create error sink");
None
}
}));
let file_sink_cfg = config.file_sink.clone().unwrap_or_default();
let (handles, shutdown_txs) = Self::start_workers(WorkerParams {
config: config.clone(),
receiver,
console_receiver,
control_rx,
control_tx: control_tx.clone(),
metrics: metrics.clone(),
console_sink: console_sink.clone(),
error_sink: error_sink.clone(),
effective_capacity: effective_capacity.clone(),
file_sink_factory: Box::new(move || {
FileSink::new(file_sink_cfg.clone()).map(|s| Box::new(s) as Box<dyn LogSink>)
}),
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
db_sink_factory: Box::new(
|db: Arc<dyn crate::integrations::Database>, metrics: Arc<Metrics>| {
let sink = crate::support::io::DatabaseSink::new(db)?;
let rt = tokio::runtime::Handle::current();
rt.block_on(async { sink.set_metrics(metrics).await });
Ok(Box::new(sink) as Box<dyn LogSink>)
},
),
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
database,
})?;
let manager = Self {
config,
sender,
console_sender,
shutdown_txs,
console_sink,
metrics,
worker_handles: Mutex::new(handles),
control_tx,
effective_capacity: effective_capacity.clone(),
#[cfg(feature = "http")]
http_server_handle: Mutex::new(None),
cache: None,
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
database: None,
};
Ok((manager, subscriber, filter))
}
pub fn builder() -> LoggerBuilder {
LoggerBuilder::default()
}
pub async fn from_file<P: AsRef<Path>>(path: P) -> Result<Self, InklogError> {
let content = std::fs::read_to_string(path.as_ref()).map_err(|e| {
let mut args = fluent_bundle::FluentArgs::new();
args.set("err", e.to_string());
InklogError::ConfigError(crate::i18n::tr_args("config-failed_read_config", args))
})?;
let config: InklogConfig = toml::from_str(&content).map_err(|e| {
let mut args = fluent_bundle::FluentArgs::new();
args.set("err", e.to_string());
InklogError::ConfigError(crate::i18n::tr_args("config-failed_parse_config", args))
})?;
Self::with_config(config).await
}
pub async fn load() -> Result<Self, InklogError> {
let config = InklogConfig::load_sync().map_err(|e| {
let mut args = fluent_bundle::FluentArgs::new();
args.set("err", e.to_string());
InklogError::ConfigError(crate::i18n::tr_args("config-failed_load_config", args))
})?;
Self::with_config(config).await
}
pub fn get_health_status(&self) -> HealthStatus {
let channel_len = self.sender.len();
let channel_cap = self.effective_capacity.load(Ordering::Acquire);
self.metrics.get_status(channel_len, channel_cap)
}
pub fn recover_sink(&self, sink_name: &str) -> Result<(), InklogError> {
self.control_tx
.send(SinkControlMessage::RecoverSink(sink_name.to_string()))
.map_err(|e| {
let mut args = fluent_bundle::FluentArgs::new();
args.set("err", e.to_string());
InklogError::ChannelError(crate::i18n::tr_args("config-failed_send_recovery", args))
})
}
pub fn effective_channel_capacity(&self) -> usize {
self.effective_capacity.load(Ordering::Acquire)
}
pub fn channel_len(&self) -> usize {
self.sender.len()
}
pub fn trigger_recovery_for_unhealthy_sinks(&self) -> Result<Vec<String>, InklogError> {
let health_status = self.get_health_status();
let mut recovered_sinks = Vec::new();
for (sink_name, sink_status) in &health_status.sinks {
if !sink_status.status.is_operational() && self.recover_sink(sink_name).is_ok() {
recovered_sinks.push(sink_name.clone());
}
}
Ok(recovered_sinks)
}
pub fn shutdown(&self) -> Result<(), InklogError> {
for tx in &self.shutdown_txs {
if tx.send(()).is_err() {
tracing::warn!("{}", crate::i18n::tr("config-shutdown_signal_lost"));
}
}
#[cfg(feature = "http")]
{
if let Ok(mut handle_guard) = self.http_server_handle.lock()
&& let Some(handle) = handle_guard.take()
{
handle.abort();
info!("HTTP server shutdown signal sent");
}
}
let handles = match self.worker_handles.lock() {
Ok(mut guard) => std::mem::take(&mut *guard),
Err(e) => {
error!("Worker handles lock poisoned: {}", e);
Vec::new()
}
};
for handle in handles {
let start = Instant::now();
while start.elapsed() < Duration::from_secs(5) {
if handle.is_finished() {
break;
}
std::thread::sleep(Duration::from_millis(10));
}
if !handle.is_finished() {
handle.abort();
}
}
Ok(())
}
}
impl Drop for LoggerManager {
fn drop(&mut self) {
let _ = self.shutdown();
}
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::Utc;
#[test]
fn test_builder_new_returns_default() {
let builder = LoggerBuilder::new();
assert_eq!(builder.config.global.level, "info");
assert!(builder.deps.cache.is_none());
assert!(builder.deps.config.is_none());
}
#[test]
fn test_builder_level_sets_config() {
let builder = LoggerBuilder::new().level("debug");
assert_eq!(builder.config.global.level, "debug");
}
#[test]
fn test_builder_level_chained() {
let builder = LoggerBuilder::new().level("trace").level("error");
assert_eq!(builder.config.global.level, "error");
}
#[test]
fn test_builder_format_sets_config() {
let builder = LoggerBuilder::new().format("{level} {message}");
assert_eq!(builder.config.global.format, "{level} {message}");
}
#[test]
fn test_builder_console_enabled_creates_config() {
let builder = LoggerBuilder::new().console(true);
assert!(builder.config.console_sink.is_some());
assert!(builder.config.console_sink.as_ref().unwrap().enabled);
}
#[test]
fn test_builder_console_disabled_keeps_some_but_disabled() {
let builder = LoggerBuilder::new().console(false);
let console = builder
.config
.console_sink
.as_ref()
.expect("console_sink should remain Some after console(false)");
assert!(!console.enabled, "console.enabled should be false");
}
#[test]
fn test_builder_file_sets_path() {
let builder = LoggerBuilder::new().file("logs/test.log");
let file_sink = builder
.config
.file_sink
.as_ref()
.expect("file_sink should be set");
assert!(file_sink.enabled);
assert_eq!(file_sink.path, std::path::PathBuf::from("logs/test.log"));
}
#[test]
fn test_builder_channel_capacity_sets_config() {
let builder = LoggerBuilder::new().channel_capacity(5000);
assert_eq!(builder.config.performance.channel_capacity, 5000);
}
#[test]
fn test_builder_worker_threads_sets_config() {
let builder = LoggerBuilder::new().worker_threads(8);
assert_eq!(builder.config.performance.worker_threads, 8);
}
#[test]
fn test_builder_console_colored_sets_config() {
let builder = LoggerBuilder::new().console(true).console_colored(false);
assert!(!builder.config.console_sink.as_ref().unwrap().colored);
}
#[test]
fn test_builder_file_max_size_sets_config() {
let builder = LoggerBuilder::new()
.file("logs/test.log")
.file_max_size("50MB");
assert_eq!(builder.config.file_sink.as_ref().unwrap().max_size, "50MB");
}
#[test]
fn test_builder_file_compress_sets_config() {
let builder = LoggerBuilder::new()
.file("logs/test.log")
.file_compress(false);
assert!(!builder.config.file_sink.as_ref().unwrap().compress);
}
#[test]
fn test_builder_file_rotation_time_sets_config() {
let builder = LoggerBuilder::new()
.file("logs/test.log")
.file_rotation_time("hourly");
assert_eq!(
builder.config.file_sink.as_ref().unwrap().rotation_time,
"hourly"
);
}
#[test]
fn test_builder_file_keep_files_sets_config() {
let builder = LoggerBuilder::new()
.file("logs/test.log")
.file_keep_files(7);
assert_eq!(builder.config.file_sink.as_ref().unwrap().keep_files, 7);
}
#[cfg(feature = "http")]
#[test]
fn test_builder_enable_http_server_creates_config() {
let builder = LoggerBuilder::new().enable_http_server(true);
assert!(builder.config.http_server.is_some());
assert!(builder.config.http_server.as_ref().unwrap().enabled);
}
#[cfg(feature = "http")]
#[test]
fn test_builder_http_host_sets_config() {
let builder = LoggerBuilder::new()
.enable_http_server(true)
.http_host("0.0.0.0");
assert_eq!(builder.config.http_server.as_ref().unwrap().host, "0.0.0.0");
}
#[cfg(feature = "http")]
#[test]
fn test_builder_http_port_sets_config() {
let builder = LoggerBuilder::new()
.enable_http_server(true)
.http_port(8080);
assert_eq!(builder.config.http_server.as_ref().unwrap().port, 8080);
}
#[test]
fn test_builder_full_chain() {
let builder = LoggerBuilder::new()
.level("warn")
.format("{message}")
.console(true)
.console_colored(false)
.file("logs/app.log")
.file_max_size("200MB")
.file_compress(true)
.file_rotation_time("hourly")
.file_keep_files(14)
.channel_capacity(20000)
.worker_threads(4);
assert_eq!(builder.config.global.level, "warn");
assert_eq!(builder.config.global.format, "{message}");
assert!(builder.config.console_sink.as_ref().unwrap().enabled);
assert!(!builder.config.console_sink.as_ref().unwrap().colored);
assert_eq!(
builder.config.file_sink.as_ref().unwrap().path,
std::path::PathBuf::from("logs/app.log")
);
assert_eq!(builder.config.file_sink.as_ref().unwrap().max_size, "200MB");
assert!(builder.config.file_sink.as_ref().unwrap().compress);
assert_eq!(
builder.config.file_sink.as_ref().unwrap().rotation_time,
"hourly"
);
assert_eq!(builder.config.file_sink.as_ref().unwrap().keep_files, 14);
assert_eq!(builder.config.performance.channel_capacity, 20000);
assert_eq!(builder.config.performance.worker_threads, 4);
}
#[test]
fn test_logger_dependencies_default_all_none() {
let deps = LoggerDependencies::default();
assert!(deps.cache.is_none());
assert!(deps.config.is_none());
}
#[test]
fn test_logger_dependencies_debug_format() {
let deps = LoggerDependencies::default();
let debug_str = format!("{:?}", deps);
assert!(debug_str.contains("cache"));
assert!(debug_str.contains("config"));
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_logger_manager_new_creates_instance() {
let manager = LoggerManager::new()
.await
.expect("Failed to create manager");
assert!(manager.effective_channel_capacity() > 0);
assert_eq!(manager.channel_len(), 0);
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_logger_manager_with_config_custom() {
let config = InklogConfig {
global: crate::GlobalConfig {
level: "debug".to_string(),
..Default::default()
},
performance: crate::PerformanceConfig {
channel_capacity: 5000,
worker_threads: 2,
..Default::default()
},
..Default::default()
};
let manager = LoggerManager::with_config(config)
.await
.expect("Failed to create manager with config");
assert_eq!(manager.effective_channel_capacity(), 5000);
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_logger_manager_get_health_status() {
let manager = LoggerManager::new()
.await
.expect("Failed to create manager");
let health = manager.get_health_status();
let _ = health;
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_logger_manager_shutdown_is_idempotent() {
let manager = LoggerManager::new()
.await
.expect("Failed to create manager");
let result1 = manager.shutdown();
assert!(result1.is_ok(), "First shutdown should succeed");
let result2 = manager.shutdown();
let _ = result2;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_logger_manager_builder_creates_working_instance() {
let manager = LoggerManager::builder()
.level("info")
.console(true)
.channel_capacity(1000)
.worker_threads(1)
.build()
.await
.expect("Failed to build manager");
assert_eq!(manager.effective_channel_capacity(), 1000);
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_logger_manager_with_dependencies_injects_cache() {
use crate::integrations::MockCache;
let deps = LoggerDependencies {
cache: Some(Arc::new(MockCache::new())),
config: None,
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
database: None,
};
let manager = LoggerManager::with_dependencies(deps)
.await
.expect("Failed to create manager with deps");
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_logger_manager_with_dependencies_injects_config() {
use crate::integrations::InklogConfigAdapter;
let config = InklogConfig::default();
let deps = LoggerDependencies {
cache: None,
config: Some(Arc::new(InklogConfigAdapter::from_config(config))),
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
database: None,
};
let manager = LoggerManager::with_dependencies(deps)
.await
.expect("Failed to create manager with config provider");
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_logger_manager_trigger_recovery_for_unhealthy_sinks() {
let manager = LoggerManager::new()
.await
.expect("Failed to create manager");
let result = manager.trigger_recovery_for_unhealthy_sinks();
assert!(result.is_ok(), "Trigger recovery should succeed");
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_logger_manager_builder_with_explicit_config() {
let manager = LoggerManager::builder()
.level("info")
.channel_capacity(2000)
.worker_threads(1)
.build()
.await
.expect("Failed to build manager");
assert_eq!(manager.effective_channel_capacity(), 2000);
let _ = manager.shutdown();
}
#[test]
fn test_builder_console_stderr_levels_with_existing_console() {
let builder = LoggerBuilder::new()
.console(true)
.console_stderr_levels(&["error", "warn"]);
let console = builder.config.console_sink.as_ref().expect("console_sink");
assert_eq!(
console.stderr_levels,
vec!["error".to_string(), "warn".to_string()]
);
}
#[test]
fn test_builder_console_stderr_levels_creates_new_when_absent() {
let mut builder = LoggerBuilder::new();
builder.config.console_sink = None;
let builder = builder.console_stderr_levels(&["error"]);
let console = builder
.config
.console_sink
.as_ref()
.expect("console_sink should be created");
assert_eq!(console.stderr_levels, vec!["error".to_string()]);
}
#[test]
fn test_builder_console_colored_true_creates_new_when_absent() {
let mut builder = LoggerBuilder::new();
builder.config.console_sink = None;
let builder = builder.console_colored(true);
let console = builder
.config
.console_sink
.as_ref()
.expect("console_sink should be created when colored=true");
assert!(console.colored);
}
#[test]
fn test_builder_file_max_size_without_file_creates_new() {
let builder = LoggerBuilder::new().file_max_size("50MB");
let file = builder
.config
.file_sink
.as_ref()
.expect("file_sink should be created");
assert_eq!(file.max_size, "50MB");
}
#[test]
fn test_builder_file_compress_without_file_creates_new() {
let builder = LoggerBuilder::new().file_compress(false);
let file = builder
.config
.file_sink
.as_ref()
.expect("file_sink should be created");
assert!(!file.compress);
}
#[test]
fn test_builder_file_rotation_time_without_file_creates_new() {
let builder = LoggerBuilder::new().file_rotation_time("daily");
let file = builder
.config
.file_sink
.as_ref()
.expect("file_sink should be created");
assert_eq!(file.rotation_time, "daily");
}
#[test]
fn test_builder_file_keep_files_without_file_creates_new() {
let builder = LoggerBuilder::new().file_keep_files(3);
let file = builder
.config
.file_sink
.as_ref()
.expect("file_sink should be created");
assert_eq!(file.keep_files, 3);
}
#[cfg(feature = "http")]
#[test]
fn test_builder_http_host_without_enable_creates_new() {
let builder = LoggerBuilder::new().http_host("0.0.0.0");
let http = builder
.config
.http_server
.as_ref()
.expect("http_server should be created");
assert_eq!(http.host, "0.0.0.0");
}
#[cfg(feature = "http")]
#[test]
fn test_builder_http_port_without_enable_creates_new() {
let builder = LoggerBuilder::new().http_port(9091);
let http = builder
.config
.http_server
.as_ref()
.expect("http_server should be created");
assert_eq!(http.port, 9091);
}
#[cfg(feature = "http")]
#[test]
fn test_builder_http_metrics_path_with_existing() {
let builder = LoggerBuilder::new()
.enable_http_server(true)
.http_metrics_path("/prom");
let http = builder.config.http_server.as_ref().expect("http_server");
assert_eq!(http.metrics_path, "/prom");
}
#[cfg(feature = "http")]
#[test]
fn test_builder_http_metrics_path_creates_new() {
let builder = LoggerBuilder::new().http_metrics_path("/m");
let http = builder
.config
.http_server
.as_ref()
.expect("http_server should be created");
assert_eq!(http.metrics_path, "/m");
}
#[cfg(feature = "http")]
#[test]
fn test_builder_http_health_path_with_existing() {
let builder = LoggerBuilder::new()
.enable_http_server(true)
.http_health_path("/healthz");
let http = builder.config.http_server.as_ref().expect("http_server");
assert_eq!(http.health_path, "/healthz");
}
#[cfg(feature = "http")]
#[test]
fn test_builder_http_health_path_creates_new() {
let builder = LoggerBuilder::new().http_health_path("/h");
let http = builder
.config
.http_server
.as_ref()
.expect("http_server should be created");
assert_eq!(http.health_path, "/h");
}
#[cfg(feature = "http")]
#[test]
fn test_builder_http_error_mode_warn() {
let builder = LoggerBuilder::new()
.enable_http_server(true)
.http_error_mode("warn");
let http = builder.config.http_server.as_ref().expect("http_server");
assert!(matches!(http.error_mode, crate::HttpErrorMode::Warn));
}
#[cfg(feature = "http")]
#[test]
fn test_builder_http_error_mode_strict() {
let builder = LoggerBuilder::new()
.enable_http_server(true)
.http_error_mode("strict");
let http = builder.config.http_server.as_ref().expect("http_server");
assert!(matches!(http.error_mode, crate::HttpErrorMode::Strict));
}
#[cfg(feature = "http")]
#[test]
fn test_builder_http_error_mode_unknown_falls_back_to_default() {
let builder = LoggerBuilder::new()
.enable_http_server(true)
.http_error_mode("invalid-mode");
let http = builder.config.http_server.as_ref().expect("http_server");
assert!(matches!(http.error_mode, crate::HttpErrorMode::Strict));
}
#[cfg(feature = "http")]
#[test]
fn test_builder_http_error_mode_creates_new() {
let builder = LoggerBuilder::new().http_error_mode("warn");
let http = builder
.config
.http_server
.as_ref()
.expect("http_server should be created");
assert!(matches!(http.error_mode, crate::HttpErrorMode::Warn));
}
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
#[test]
fn test_builder_database_sets_config() {
let builder = LoggerBuilder::new().database("postgres://localhost/logs");
let db = builder
.config
.database_sink
.as_ref()
.expect("database_sink should be set");
assert!(db.enabled);
assert_eq!(db.url, "postgres://localhost/logs");
assert_eq!(db.name, "default");
}
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
#[test]
fn test_builder_with_database_injects_dep() {
use crate::integrations::MockDatabaseAdapter;
let builder = LoggerBuilder::new().with_database(Arc::new(MockDatabaseAdapter::new()));
assert!(builder.deps.database.is_some());
}
#[test]
fn test_builder_cache_injects_dep() {
use crate::integrations::MockCache;
let builder = LoggerBuilder::new().cache(Arc::new(MockCache::new()));
assert!(builder.deps.cache.is_some());
}
#[test]
fn test_builder_config_injects_dep() {
use crate::integrations::MockConfig;
let builder = LoggerBuilder::new().config(Arc::new(MockConfig::new()));
assert!(builder.deps.config.is_some());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_builder_build_with_cache_injection_mixed_mode() {
use crate::integrations::MockCache;
let manager = LoggerManager::builder()
.level("info")
.channel_capacity(1500)
.worker_threads(1)
.cache(Arc::new(MockCache::new()))
.build()
.await
.expect("Failed to build manager with cache injection");
assert_eq!(manager.effective_channel_capacity(), 1500);
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_builder_build_with_config_injection() {
use crate::integrations::MockConfig;
let manager = LoggerManager::builder()
.config(Arc::new(MockConfig::new()))
.worker_threads(1)
.build()
.await
.expect("Failed to build manager with config injection");
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_logger_manager_recover_sink_on_live_manager() {
let dir = tempfile::tempdir().expect("Failed to create tempdir");
let log_path = dir.path().join("app.log");
let manager = LoggerManager::builder()
.channel_capacity(1000)
.worker_threads(1)
.file(log_path)
.build()
.await
.expect("Failed to build manager");
let result = manager.recover_sink("file");
assert!(
result.is_ok(),
"recover_sink on live manager should succeed"
);
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_logger_manager_from_file_loads_valid_config() {
let dir = tempfile::tempdir().expect("Failed to create tempdir");
let config_path = dir.path().join("inklog_config.toml");
let toml_content = r#"
[global]
level = "debug"
[performance]
channel_capacity = 3000
worker_threads = 1
"#;
std::fs::write(&config_path, toml_content).expect("Failed to write config");
let manager = LoggerManager::from_file(&config_path)
.await
.expect("Failed to load manager from file");
assert_eq!(manager.effective_channel_capacity(), 3000);
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_logger_manager_from_file_missing_path_returns_error() {
let dir = tempfile::tempdir().expect("Failed to create tempdir");
let missing = dir.path().join("nonexistent.toml");
let result = LoggerManager::from_file(&missing).await;
assert!(result.is_err(), "from_file with missing path should error");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_logger_manager_from_file_invalid_toml_returns_error() {
let dir = tempfile::tempdir().expect("Failed to create tempdir");
let config_path = dir.path().join("invalid.toml");
std::fs::write(&config_path, "this is = = not valid toml [[[")
.expect("Failed to write config");
let result = LoggerManager::from_file(&config_path).await;
assert!(result.is_err(), "from_file with invalid toml should error");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_manager_with_config_trace_level() {
let config = InklogConfig {
global: crate::GlobalConfig {
level: "trace".to_string(),
..Default::default()
},
performance: crate::PerformanceConfig {
channel_capacity: 1000,
worker_threads: 1,
..Default::default()
},
..Default::default()
};
let manager = LoggerManager::with_config(config)
.await
.expect("Failed to create manager with trace level");
assert_eq!(manager.effective_channel_capacity(), 1000);
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_manager_with_config_warn_level() {
let config = InklogConfig {
global: crate::GlobalConfig {
level: "warn".to_string(),
..Default::default()
},
performance: crate::PerformanceConfig {
channel_capacity: 1000,
worker_threads: 1,
..Default::default()
},
..Default::default()
};
let manager = LoggerManager::with_config(config)
.await
.expect("Failed to create manager with warn level");
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_manager_with_config_error_level() {
let config = InklogConfig {
global: crate::GlobalConfig {
level: "error".to_string(),
..Default::default()
},
performance: crate::PerformanceConfig {
channel_capacity: 1000,
worker_threads: 1,
..Default::default()
},
..Default::default()
};
let manager = LoggerManager::with_config(config)
.await
.expect("Failed to create manager with error level");
let _ = manager.shutdown();
}
#[cfg(feature = "http")]
#[test]
fn test_builder_enable_http_server_false_when_exists() {
let builder = LoggerBuilder::new()
.enable_http_server(true)
.enable_http_server(false);
let http = builder
.config
.http_server
.as_ref()
.expect("http_server should exist");
assert!(!http.enabled, "http.enabled should be false after disable");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_manager_file_sink_writes_record_to_file() {
let dir = tempfile::tempdir().expect("Failed to create tempdir");
let log_path = dir.path().join("worker_test.log");
let manager = LoggerManager::builder()
.channel_capacity(500)
.worker_threads(1)
.file(&log_path)
.build()
.await
.expect("Failed to build manager with file sink");
let record = Arc::new(LogRecord {
timestamp: Utc::now(),
level: "INFO".to_string(),
target: "worker_test".to_string(),
message: "worker_write_unique_marker_12345".to_string(),
fields: std::collections::HashMap::new(),
file: None,
line: None,
thread_id: "test-thread".to_string(),
});
manager
.sender
.send(record)
.expect("Failed to send record to file worker");
std::thread::sleep(Duration::from_millis(300));
let _ = manager.shutdown();
let content =
std::fs::read_to_string(&log_path).expect("Log file should exist after shutdown");
assert!(
content.contains("worker_write_unique_marker_12345"),
"Log file should contain the sent message"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_manager_file_sink_drains_multiple_records_on_shutdown() {
let dir = tempfile::tempdir().expect("Failed to create tempdir");
let log_path = dir.path().join("drain_test.log");
let manager = LoggerManager::builder()
.channel_capacity(500)
.worker_threads(1)
.file(&log_path)
.build()
.await
.expect("Failed to build manager");
for i in 0..10u32 {
let record = Arc::new(LogRecord {
timestamp: Utc::now(),
level: "INFO".to_string(),
target: "drain_test".to_string(),
message: format!("drain_record_{:02}", i),
fields: std::collections::HashMap::new(),
file: None,
line: None,
thread_id: "test-thread".to_string(),
});
manager.sender.send(record).expect("Failed to send record");
}
let _ = manager.shutdown();
let content = std::fs::read_to_string(&log_path).expect("Log file should exist");
for i in 0..10u32 {
let marker = format!("drain_record_{:02}", i);
assert!(
content.contains(&marker),
"Log file should contain '{}'",
marker
);
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_recover_sink_multiple_commands_to_live_manager() {
let dir = tempfile::tempdir().expect("Failed to create tempdir");
let log_path = dir.path().join("recover_test.log");
let manager = LoggerManager::builder()
.channel_capacity(500)
.worker_threads(1)
.file(&log_path)
.build()
.await
.expect("Failed to build manager");
let r1 = manager.recover_sink("file");
let r2 = manager.recover_sink("database");
let r3 = manager.recover_sink("unknown_sink");
assert!(r1.is_ok(), "recover_sink('file') should succeed");
assert!(r2.is_ok(), "recover_sink('database') should succeed");
assert!(r3.is_ok(), "recover_sink('unknown') should succeed");
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_manager_console_sink_processes_record() {
let manager = LoggerManager::builder()
.channel_capacity(500)
.worker_threads(1)
.console(true)
.build()
.await
.expect("Failed to build manager");
let record = Arc::new(LogRecord {
timestamp: Utc::now(),
level: "INFO".to_string(),
target: "console_test".to_string(),
message: "console_marker_98765".to_string(),
fields: std::collections::HashMap::new(),
file: None,
line: None,
thread_id: "test-thread".to_string(),
});
manager
.console_sender
.send(record)
.expect("Failed to send record to console worker");
std::thread::sleep(Duration::from_millis(200));
let _ = manager.shutdown();
}
#[cfg(feature = "http")]
fn find_available_http_port() -> u16 {
let listener = std::net::TcpListener::bind("127.0.0.1:0")
.expect("Failed to bind to find available port");
let port = listener
.local_addr()
.expect("Failed to get local addr")
.port();
drop(listener);
port
}
#[cfg(feature = "http")]
async fn wait_for_http_server(host: &str, port: u16) -> bool {
let url = format!("http://{}:{}", host, port);
for _ in 0..80 {
if reqwest::get(&url).await.is_ok() {
return true;
}
tokio::time::sleep(std::time::Duration::from_millis(25)).await;
}
false
}
#[cfg(feature = "http")]
fn http_test_config(port: u16) -> InklogConfig {
InklogConfig {
http_server: Some(crate::HttpServerConfig {
enabled: true,
host: "127.0.0.1".to_string(),
port,
error_mode: crate::HttpErrorMode::Warn,
..Default::default()
}),
performance: crate::PerformanceConfig {
channel_capacity: 1000,
worker_threads: 1,
..Default::default()
},
..Default::default()
}
}
#[cfg(feature = "http")]
fn http_test_config_with_auth(port: u16, token_env: &str) -> InklogConfig {
let mut config = http_test_config(port);
let http = config
.http_server
.as_mut()
.expect("http_server should be set");
http.auth = Some(crate::HttpAuthConfig {
enabled: true,
token_env: token_env.to_string(),
});
config
}
#[cfg(feature = "http")]
fn http_test_config_with_whitelist(port: u16, whitelist: Vec<String>) -> InklogConfig {
let mut config = http_test_config(port);
let http = config
.http_server
.as_mut()
.expect("http_server should be set");
http.ip_whitelist = Some(whitelist);
config
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn test_with_config_http_warn_mode_continues_on_startup_error() {
let config = InklogConfig {
http_server: Some(crate::HttpServerConfig {
enabled: true,
host: "invalid host with spaces".to_string(),
port: 9090,
error_mode: crate::HttpErrorMode::Warn,
..Default::default()
}),
performance: crate::PerformanceConfig {
channel_capacity: 1000,
worker_threads: 1,
..Default::default()
},
..Default::default()
};
let manager = LoggerManager::with_config(config)
.await
.expect("Warn mode should return Ok despite HTTP server startup error");
let _ = manager.shutdown();
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn test_with_config_http_strict_mode_returns_error_on_invalid_host() {
let config = InklogConfig {
http_server: Some(crate::HttpServerConfig {
enabled: true,
host: "invalid host with spaces".to_string(),
port: 9091,
error_mode: crate::HttpErrorMode::Strict,
..Default::default()
}),
performance: crate::PerformanceConfig {
channel_capacity: 1000,
worker_threads: 1,
..Default::default()
},
..Default::default()
};
match LoggerManager::with_config(config).await {
Err(InklogError::ConfigError(msg)) => {
assert!(
msg.contains("Invalid HTTP server address"),
"Error should mention invalid HTTP server address, got: {}",
msg
);
}
Err(other) => panic!("Expected ConfigError, got {:?}", other),
Ok(_) => panic!("Strict mode should return Err on invalid HTTP address"),
}
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn test_http_server_health_endpoint_returns_json() {
let port = find_available_http_port();
let manager = LoggerManager::with_config(http_test_config(port))
.await
.expect("Manager should start with HTTP server");
assert!(
wait_for_http_server("127.0.0.1", port).await,
"HTTP server should become reachable on port {}",
port
);
let resp = reqwest::get(format!("http://127.0.0.1:{}/health", port))
.await
.expect("GET /health should succeed");
assert_eq!(
resp.status(),
reqwest::StatusCode::OK,
"health endpoint should return 200"
);
let body: serde_json::Value = resp.json().await.expect("body should be JSON");
assert!(body.is_object(), "health response should be a JSON object");
assert!(
body.get("overall_status").is_some(),
"health response should contain overall_status field"
);
let _ = manager.shutdown();
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn test_http_server_metrics_endpoint_returns_prometheus() {
let port = find_available_http_port();
let manager = LoggerManager::with_config(http_test_config(port))
.await
.expect("Manager should start with HTTP server");
assert!(
wait_for_http_server("127.0.0.1", port).await,
"HTTP server should become reachable on port {}",
port
);
let resp = reqwest::get(format!("http://127.0.0.1:{}/metrics", port))
.await
.expect("GET /metrics should succeed");
assert_eq!(
resp.status(),
reqwest::StatusCode::OK,
"metrics endpoint should return 200"
);
let body = resp.text().await.expect("body should be text");
assert!(
body.contains("# HELP") && body.contains("inklog_"),
"metrics response should be in Prometheus format, got: {}",
body
);
let _ = manager.shutdown();
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn test_http_server_custom_paths_work() {
let port = find_available_http_port();
let mut config = http_test_config(port);
{
let http = config
.http_server
.as_mut()
.expect("http_server should be set");
http.health_path = "/custom-health".to_string();
http.metrics_path = "/custom-metrics".to_string();
}
let manager = LoggerManager::with_config(config)
.await
.expect("Manager should start with HTTP server");
assert!(
wait_for_http_server("127.0.0.1", port).await,
"HTTP server should become reachable"
);
let resp = reqwest::get(format!("http://127.0.0.1:{}/custom-health", port))
.await
.expect("GET /custom-health should succeed");
assert_eq!(
resp.status(),
reqwest::StatusCode::OK,
"custom health path should return 200"
);
let resp = reqwest::get(format!("http://127.0.0.1:{}/custom-metrics", port))
.await
.expect("GET /custom-metrics should succeed");
assert_eq!(
resp.status(),
reqwest::StatusCode::OK,
"custom metrics path should return 200"
);
let resp = reqwest::get(format!("http://127.0.0.1:{}/health", port))
.await
.expect("GET /health should succeed");
assert_eq!(
resp.status(),
reqwest::StatusCode::NOT_FOUND,
"default health path should return 404 when customized"
);
let _ = manager.shutdown();
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn test_http_server_auth_disabled_allows_access_without_header() {
let port = find_available_http_port();
let manager = LoggerManager::with_config(http_test_config(port))
.await
.expect("Manager should start with HTTP server");
assert!(
wait_for_http_server("127.0.0.1", port).await,
"HTTP server should become reachable"
);
let resp = reqwest::get(format!("http://127.0.0.1:{}/health", port))
.await
.expect("GET /health should succeed");
assert_eq!(
resp.status(),
reqwest::StatusCode::OK,
"auth disabled should allow access without Authorization header"
);
let _ = manager.shutdown();
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn test_http_server_auth_missing_token_env_fails_to_start() {
let port = find_available_http_port();
let token_env = "INKLOG_TEST_TOKEN_MISSING_ENV_VAR";
unsafe {
std::env::remove_var(token_env);
}
let mut config = http_test_config_with_auth(port, token_env);
config.http_server.as_mut().unwrap().error_mode = crate::HttpErrorMode::Strict;
let result = LoggerManager::with_config(config).await;
let err_msg = match result {
Err(e) => format!("{}", e),
Ok(_) => panic!("vuln-0003: missing token env should fail to start (fail-closed)"),
};
assert!(
err_msg.contains("token env var") && err_msg.contains("is not set"),
"error should explain token env misconfiguration, got: {}",
err_msg
);
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn test_http_server_auth_empty_token_env_fails_to_start() {
let port = find_available_http_port();
let token_env = "INKLOG_TEST_TOKEN_EMPTY_ENV_VAR";
unsafe {
std::env::set_var(token_env, "");
}
let mut config = http_test_config_with_auth(port, token_env);
config.http_server.as_mut().unwrap().error_mode = crate::HttpErrorMode::Strict;
let result = LoggerManager::with_config(config).await;
let err_msg = match result {
Err(e) => format!("{}", e),
Ok(_) => panic!("vuln-0003: empty token env should fail to start (fail-closed)"),
};
assert!(
err_msg.contains("is empty"),
"error should explain token env is empty, got: {}",
err_msg
);
unsafe {
std::env::remove_var(token_env);
}
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn test_vuln_0003_env_var_change_after_start_does_not_affect_auth() {
let port = find_available_http_port();
let token_env = "INKLOG_TEST_TOKEN_VULN_0003";
let original_token = "original-secret-vuln-0003";
unsafe {
std::env::set_var(token_env, original_token);
}
let manager = LoggerManager::with_config(http_test_config_with_auth(port, token_env))
.await
.expect("Manager should start with valid token");
assert!(
wait_for_http_server("127.0.0.1", port).await,
"HTTP server should become reachable"
);
let client = reqwest::Client::builder()
.build()
.expect("Failed to build reqwest client");
let resp = client
.get(format!("http://127.0.0.1:{}/health", port))
.bearer_auth(original_token)
.send()
.await
.expect("Request with original token should succeed");
assert_eq!(
resp.status(),
reqwest::StatusCode::OK,
"original token should work"
);
let tampered_token = "tampered-by-attacker";
unsafe {
std::env::set_var(token_env, tampered_token);
}
let resp = client
.get(format!("http://127.0.0.1:{}/health", port))
.bearer_auth(tampered_token)
.send()
.await
.expect("Request with tampered token should still get a response");
assert_eq!(
resp.status(),
reqwest::StatusCode::UNAUTHORIZED,
"vuln-0003: tampered env var token should NOT work (cached token at startup wins)"
);
let resp = client
.get(format!("http://127.0.0.1:{}/health", port))
.bearer_auth(original_token)
.send()
.await
.expect("Request with original token should still succeed");
assert_eq!(
resp.status(),
reqwest::StatusCode::OK,
"vuln-0003: original token should still work after env var tampering"
);
let _ = manager.shutdown();
unsafe {
std::env::remove_var(token_env);
}
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn test_http_server_auth_valid_token_returns_200() {
let port = find_available_http_port();
let token_env = "INKLOG_TEST_TOKEN_VALID";
let token_value = "secret-token-12345";
unsafe {
std::env::set_var(token_env, token_value);
}
let manager = LoggerManager::with_config(http_test_config_with_auth(port, token_env))
.await
.expect("Manager should start with HTTP server");
assert!(
wait_for_http_server("127.0.0.1", port).await,
"HTTP server should become reachable"
);
let client = reqwest::Client::builder()
.build()
.expect("Failed to build reqwest client");
let resp = client
.get(format!("http://127.0.0.1:{}/health", port))
.bearer_auth(token_value)
.send()
.await
.expect("Request with valid token should succeed");
assert_eq!(
resp.status(),
reqwest::StatusCode::OK,
"valid Bearer token should return 200"
);
let _ = manager.shutdown();
unsafe {
std::env::remove_var(token_env);
}
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn test_http_server_auth_invalid_token_returns_401() {
let port = find_available_http_port();
let token_env = "INKLOG_TEST_TOKEN_INVALID";
unsafe {
std::env::set_var(token_env, "correct-secret");
}
let manager = LoggerManager::with_config(http_test_config_with_auth(port, token_env))
.await
.expect("Manager should start with HTTP server");
assert!(
wait_for_http_server("127.0.0.1", port).await,
"HTTP server should become reachable"
);
let client = reqwest::Client::builder()
.build()
.expect("Failed to build reqwest client");
let resp = client
.get(format!("http://127.0.0.1:{}/health", port))
.bearer_auth("wrong-secret")
.send()
.await
.expect("Request with invalid token should still get a response");
assert_eq!(
resp.status(),
reqwest::StatusCode::UNAUTHORIZED,
"invalid Bearer token should return 401"
);
let body = resp.text().await.expect("body should be text");
assert!(
body.contains("Invalid token"),
"response should indicate invalid token, got: {}",
body
);
let _ = manager.shutdown();
unsafe {
std::env::remove_var(token_env);
}
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn test_http_server_auth_missing_header_returns_401() {
let port = find_available_http_port();
let token_env = "INKLOG_TEST_TOKEN_MISSING_HEADER";
unsafe {
std::env::set_var(token_env, "some-secret");
}
let manager = LoggerManager::with_config(http_test_config_with_auth(port, token_env))
.await
.expect("Manager should start with HTTP server");
assert!(
wait_for_http_server("127.0.0.1", port).await,
"HTTP server should become reachable"
);
let resp = reqwest::get(format!("http://127.0.0.1:{}/health", port))
.await
.expect("Request without header should still get a response");
assert_eq!(
resp.status(),
reqwest::StatusCode::UNAUTHORIZED,
"missing Authorization header should return 401"
);
let body = resp.text().await.expect("body should be text");
assert!(
body.contains("Missing or invalid Authorization header"),
"response should indicate missing header, got: {}",
body
);
let _ = manager.shutdown();
unsafe {
std::env::remove_var(token_env);
}
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn test_http_server_ip_whitelist_allows_exact_match() {
let port = find_available_http_port();
let config = http_test_config_with_whitelist(port, vec!["127.0.0.1".to_string()]);
let manager = LoggerManager::with_config(config)
.await
.expect("Manager should start with HTTP server");
assert!(
wait_for_http_server("127.0.0.1", port).await,
"HTTP server should become reachable"
);
let resp = reqwest::get(format!("http://127.0.0.1:{}/health", port))
.await
.expect("GET /health should succeed");
assert_eq!(
resp.status(),
reqwest::StatusCode::OK,
"exact IP match in whitelist should allow access"
);
let _ = manager.shutdown();
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn test_http_server_ip_whitelist_rejects_non_match() {
let port = find_available_http_port();
let config = http_test_config_with_whitelist(port, vec!["10.0.0.1".to_string()]);
let manager = LoggerManager::with_config(config)
.await
.expect("Manager should start with HTTP server");
assert!(
wait_for_http_server("127.0.0.1", port).await,
"HTTP server should become reachable"
);
let resp = reqwest::get(format!("http://127.0.0.1:{}/health", port))
.await
.expect("GET /health should still get a response");
assert_eq!(
resp.status(),
reqwest::StatusCode::FORBIDDEN,
"non-matching IP should be forbidden"
);
let body = resp.text().await.expect("body should be text");
assert!(
body.contains("IP not in whitelist"),
"response should indicate IP rejection, got: {}",
body
);
let _ = manager.shutdown();
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn test_http_server_ip_whitelist_allows_wildcard() {
let port = find_available_http_port();
let config = http_test_config_with_whitelist(port, vec!["127.0.*".to_string()]);
let manager = LoggerManager::with_config(config)
.await
.expect("Manager should start with HTTP server");
assert!(
wait_for_http_server("127.0.0.1", port).await,
"HTTP server should become reachable"
);
let resp = reqwest::get(format!("http://127.0.0.1:{}/health", port))
.await
.expect("GET /health should succeed");
assert_eq!(
resp.status(),
reqwest::StatusCode::OK,
"wildcard 127.0.* should match 127.0.0.1"
);
let _ = manager.shutdown();
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn test_http_server_ip_whitelist_allows_cidr() {
let port = find_available_http_port();
let config = http_test_config_with_whitelist(port, vec!["127.0.0.0/8".to_string()]);
let manager = LoggerManager::with_config(config)
.await
.expect("Manager should start with HTTP server");
assert!(
wait_for_http_server("127.0.0.1", port).await,
"HTTP server should become reachable"
);
let resp = reqwest::get(format!("http://127.0.0.1:{}/health", port))
.await
.expect("GET /health should succeed");
assert_eq!(
resp.status(),
reqwest::StatusCode::OK,
"CIDR 127.0.0.0/8 should contain 127.0.0.1"
);
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_build_with_deps_applies_global_config_from_provider() {
use crate::integrations::MockConfig;
let mock_config = MockConfig::new()
.with_value("global.level", "debug")
.with_value("global.format", "{level} {message}")
.with_value("global.masking_enabled", "true")
.with_value("global.auto_fallback", "true");
let deps = LoggerDependencies {
cache: None,
config: Some(Arc::new(mock_config)),
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
database: None,
};
let manager = LoggerManager::with_dependencies(deps)
.await
.expect("Failed to create manager with config provider");
let config = &manager.config;
assert_eq!(config.global.level, "debug");
assert_eq!(config.global.format, "{level} {message}");
assert!(config.global.masking_enabled);
assert!(config.global.auto_fallback);
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_build_with_deps_configures_file_sink_from_provider() {
use crate::integrations::MockConfig;
let dir = tempfile::tempdir().expect("Failed to create tempdir");
let log_path = dir.path().join("from_provider.log");
let path_str = log_path
.to_str()
.expect("path should be valid utf-8")
.to_string();
let mock_config = MockConfig::new()
.with_value("file_sink.enabled", "true")
.with_value("file_sink.path", &path_str)
.with_value("file_sink.max_size", "50MB")
.with_value("file_sink.compress", "false");
let deps = LoggerDependencies {
cache: None,
config: Some(Arc::new(mock_config)),
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
database: None,
};
let manager = LoggerManager::with_dependencies(deps)
.await
.expect("Failed to create manager with file_sink config");
let config = &manager.config;
let file_sink = config
.file_sink
.as_ref()
.expect("file_sink should be configured from provider");
assert!(file_sink.enabled);
assert_eq!(file_sink.path, std::path::PathBuf::from(&path_str));
assert_eq!(file_sink.max_size, "50MB");
assert!(!file_sink.compress);
let record = Arc::new(LogRecord {
timestamp: Utc::now(),
level: "INFO".to_string(),
target: "config_provider_test".to_string(),
message: "from_provider_unique_marker_abc123".to_string(),
fields: std::collections::HashMap::new(),
file: None,
line: None,
thread_id: "test".to_string(),
});
manager
.sender
.send(record)
.expect("Failed to send record to file worker");
std::thread::sleep(Duration::from_millis(300));
let _ = manager.shutdown();
let content =
std::fs::read_to_string(&log_path).expect("Log file should exist after write");
assert!(
content.contains("from_provider_unique_marker_abc123"),
"Log file should contain the message sent via config_provider-configured file sink"
);
}
#[cfg(feature = "http")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_build_with_deps_configures_http_server_from_provider() {
use crate::integrations::MockConfig;
let mock_config = MockConfig::new()
.with_value("http_server.enabled", "true")
.with_value("http_server.host", "127.0.0.1")
.with_value("http_server.port", "9090");
let deps = LoggerDependencies {
cache: None,
config: Some(Arc::new(mock_config)),
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
database: None,
};
let manager = LoggerManager::with_dependencies(deps)
.await
.expect("Failed to create manager with http_server config");
let config = &manager.config;
let http = config
.http_server
.as_ref()
.expect("http_server should be configured from provider");
assert!(http.enabled);
assert_eq!(http.host, "127.0.0.1");
assert_eq!(http.port, 9090);
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_build_with_deps_configures_performance_from_provider() {
use crate::integrations::MockConfig;
let mock_config = MockConfig::new()
.with_value("performance.worker_threads", "2")
.with_value("performance.channel_capacity", "3000");
let deps = LoggerDependencies {
cache: None,
config: Some(Arc::new(mock_config)),
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
database: None,
};
let manager = LoggerManager::with_dependencies(deps)
.await
.expect("Failed to create manager with performance config");
assert_eq!(
manager.effective_channel_capacity(),
3000,
"channel_capacity from config provider should be applied"
);
let config = &manager.config;
assert_eq!(config.performance.worker_threads, 2);
let _ = manager.shutdown();
}
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_build_with_deps_injects_database() {
use crate::integrations::MockDatabaseAdapter;
let deps = LoggerDependencies {
cache: None,
config: None,
database: Some(Arc::new(MockDatabaseAdapter::new())),
};
let manager = LoggerManager::with_dependencies(deps)
.await
.expect("Failed to create manager with database injection");
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_build_detached_returns_valid_components() {
let config = InklogConfig {
global: crate::GlobalConfig {
level: "warn".to_string(),
..Default::default()
},
performance: crate::PerformanceConfig {
channel_capacity: 2000,
worker_threads: 1,
..Default::default()
},
..Default::default()
};
let (manager, _subscriber, filter) = LoggerManager::build_detached(
config,
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
None,
)
.await
.expect("build_detached should succeed with valid config");
assert_eq!(
manager.effective_channel_capacity(),
2000,
"effective_channel_capacity should match config"
);
assert_eq!(
manager.channel_len(),
0,
"channel_len should be 0 for fresh manager"
);
let filter_str = filter.to_string();
assert!(
filter_str.contains("warn"),
"EnvFilter should contain config level 'warn', got: {}",
filter_str
);
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_build_detached_invalid_level_falls_back_to_info() {
let config = InklogConfig {
global: crate::GlobalConfig {
level: "invalid_level".to_string(),
..Default::default()
},
performance: crate::PerformanceConfig {
channel_capacity: 1000,
worker_threads: 1,
..Default::default()
},
..Default::default()
};
let (manager, _subscriber, filter) = LoggerManager::build_detached(
config,
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
None,
)
.await
.expect("build_detached should succeed even with invalid level");
let filter_str = filter.to_string();
assert!(
filter_str.contains("info"),
"EnvFilter should fall back to 'info' for invalid level, got: {}",
filter_str
);
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_file_worker_skips_when_file_sink_new_fails() {
let config = InklogConfig {
file_sink: Some(FileSinkConfig {
enabled: true,
path: PathBuf::from("/dev/null/subdir/file.log"),
..Default::default()
}),
performance: crate::PerformanceConfig {
channel_capacity: 1000,
worker_threads: 1,
..Default::default()
},
..Default::default()
};
let manager = LoggerManager::with_config(config)
.await
.expect("Manager should be created even if FileSink::new fails in worker");
assert_eq!(manager.effective_channel_capacity(), 1000);
let result = manager.shutdown();
assert!(result.is_ok(), "shutdown should succeed");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_file_worker_recovers_after_recover_sink_command() {
let dir = tempfile::tempdir().expect("Failed to create tempdir");
let log_path = dir.path().join("recover_worker.log");
let manager = LoggerManager::builder()
.channel_capacity(500)
.worker_threads(1)
.file(&log_path)
.build()
.await
.expect("Failed to build manager");
let record = Arc::new(LogRecord {
timestamp: Utc::now(),
level: "INFO".to_string(),
target: "recover_test".to_string(),
message: "before_recover_marker".to_string(),
fields: std::collections::HashMap::new(),
file: None,
line: None,
thread_id: "test".to_string(),
});
manager.sender.send(record).expect("Failed to send record");
std::thread::sleep(Duration::from_millis(200));
let result = manager.recover_sink("file");
assert!(
result.is_ok(),
"recover_sink('file') should succeed on live manager"
);
std::thread::sleep(Duration::from_millis(200));
let record2 = Arc::new(LogRecord {
timestamp: Utc::now(),
level: "INFO".to_string(),
target: "recover_test".to_string(),
message: "after_recover_marker".to_string(),
fields: std::collections::HashMap::new(),
file: None,
line: None,
thread_id: "test".to_string(),
});
manager.sender.send(record2).expect("Failed to send record");
std::thread::sleep(Duration::from_millis(300));
let _ = manager.shutdown();
let content =
std::fs::read_to_string(&log_path).expect("Log file should exist after recover");
assert!(
content.contains("before_recover_marker"),
"Log file should contain record sent before recover"
);
assert!(
content.contains("after_recover_marker"),
"Log file should contain record sent after recover"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_recover_sink_returns_error_when_control_channel_full() {
let dir = tempfile::tempdir().expect("Failed to create tempdir");
let log_path = dir.path().join("channel_full.log");
let manager = LoggerManager::builder()
.channel_capacity(500)
.worker_threads(1)
.file(&log_path)
.build()
.await
.expect("Failed to build manager");
let mut ok_count = 0;
let mut err_count = 0;
for _ in 0..20 {
match manager.recover_sink("file") {
Ok(_) => ok_count += 1,
Err(InklogError::ChannelError(_)) => err_count += 1,
Err(other) => panic!("Unexpected error type: {:?}", other),
}
}
assert!(
ok_count > 0,
"At least some recover_sink commands should succeed"
);
assert_eq!(ok_count + err_count, 20);
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_build_detached_creates_error_sink() {
let config = InklogConfig {
performance: crate::PerformanceConfig {
channel_capacity: 1000,
worker_threads: 1,
..Default::default()
},
..Default::default()
};
let (manager, _subscriber, _filter) = LoggerManager::build_detached(
config,
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
None,
)
.await
.expect("build_detached should succeed");
assert!(manager.effective_channel_capacity() > 0);
let _ = manager.shutdown();
}
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
#[test]
fn test_logger_dependencies_debug_includes_database_field() {
use crate::integrations::{MockCache, MockDatabaseAdapter};
let deps = LoggerDependencies {
cache: Some(Arc::new(MockCache::new())),
config: None,
database: Some(Arc::new(MockDatabaseAdapter::new())),
};
let debug_str = format!("{:?}", deps);
assert!(
debug_str.contains("cache"),
"debug should include cache field"
);
assert!(
debug_str.contains("config"),
"debug should include config field"
);
assert!(
debug_str.contains("database"),
"debug should include database field when dbnexus feature enabled"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_build_with_deps_injects_both_cache_and_config() {
use crate::integrations::{InklogConfigAdapter, MockCache};
let config = InklogConfig::default();
let deps = LoggerDependencies {
cache: Some(Arc::new(MockCache::new())),
config: Some(Arc::new(InklogConfigAdapter::from_config(config))),
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
database: None,
};
let manager = LoggerManager::with_dependencies(deps)
.await
.expect("Failed to create manager with cache and config");
let _ = manager.shutdown();
}
#[cfg(any(
feature = "sqlite",
feature = "postgres",
feature = "mysql",
feature = "duckdb"
))]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_build_with_deps_injects_all_three_deps() {
use crate::integrations::{InklogConfigAdapter, MockCache, MockDatabaseAdapter};
let config = InklogConfig::default();
let deps = LoggerDependencies {
cache: Some(Arc::new(MockCache::new())),
config: Some(Arc::new(InklogConfigAdapter::from_config(config))),
database: Some(Arc::new(MockDatabaseAdapter::new())),
};
let manager = LoggerManager::with_dependencies(deps)
.await
.expect("Failed to create manager with all deps");
let _ = manager.shutdown();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial]
async fn test_load_succeeds_with_valid_config_via_env() {
let dir = tempfile::tempdir().expect("Failed to create tempdir");
let config_path = dir.path().join("load_test.toml");
std::fs::write(&config_path, "[global]\nlevel = \"info\"\n").unwrap();
unsafe {
std::env::set_var("INKLOG_CONFIG_PATH", &config_path);
}
let manager = LoggerManager::load().await.expect("load should succeed");
let _ = manager.shutdown();
unsafe {
std::env::remove_var("INKLOG_CONFIG_PATH");
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial]
async fn test_load_returns_error_when_config_invalid_toml() {
let dir = tempfile::tempdir().expect("Failed to create tempdir");
let config_path = dir.path().join("invalid.toml");
std::fs::write(&config_path, "this is = not = valid toml\n").unwrap();
unsafe {
std::env::set_var("INKLOG_CONFIG_PATH", &config_path);
}
let result = LoggerManager::load().await;
assert!(result.is_err(), "load should fail with invalid TOML");
let err_msg = result.err().unwrap().to_string();
assert!(
err_msg.contains("Failed to load config") || err_msg.contains("Failed to parse"),
"error should mention config load failure, got: {}",
err_msg
);
unsafe {
std::env::remove_var("INKLOG_CONFIG_PATH");
}
}
#[test]
fn test_builder_console_when_none_creates_config() {
let mut builder = LoggerBuilder::new();
builder.config.console_sink = None;
let builder = builder.console(true);
let console = builder
.config
.console_sink
.as_ref()
.expect("console_sink should be Some after console(true)");
assert!(
console.enabled,
"console.enabled should be true after console(true)"
);
}
#[test]
fn test_builder_file_when_some_updates_path() {
let mut builder = LoggerBuilder::new();
builder.config.file_sink = Some(crate::FileSinkConfig::default());
let builder = builder.file("logs/updated.log");
let file_sink = builder
.config
.file_sink
.as_ref()
.expect("file_sink should remain Some");
assert!(
file_sink.enabled,
"file_sink.enabled should be true after file(path)"
);
assert_eq!(
file_sink.path,
std::path::PathBuf::from("logs/updated.log"),
"file_sink.path should be updated"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_trigger_recovery_recovers_unhealthy_sink() {
let dir = tempfile::tempdir().expect("Failed to create tempdir");
let log_path = dir.path().join("recovery_test.log");
let manager = LoggerManager::builder()
.channel_capacity(1000)
.worker_threads(1)
.file(log_path)
.build()
.await
.expect("Failed to build manager");
manager
.metrics
.update_sink_health("file", false, Some("test error".to_string()));
let result = manager.trigger_recovery_for_unhealthy_sinks();
assert!(result.is_ok(), "trigger_recovery should succeed");
let recovered = result.unwrap();
assert!(
recovered.contains(&"file".to_string()),
"recovered sinks should contain 'file', got: {:?}",
recovered
);
let _ = manager.shutdown();
}
}