use inklog::sink::LogSink;
use inklog::LoggerManager;
use serial_test::serial;
use std::time::Duration;
use tracing::{error, info};
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
#[allow(unused_imports)]
use inklog::sink::database::DatabaseSink;
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn test_e2e_logging() {
if let Ok(logger) = LoggerManager::new().await {
info!("This is an info message");
error!("This is an error message");
std::thread::sleep(Duration::from_millis(200));
logger.shutdown().expect("Failed to shutdown logger");
}
}
#[tokio::test]
async fn test_load_from_file() {
use std::io::Write;
let mut file = tempfile::NamedTempFile::new().expect("Failed to create temp file");
write!(
file,
r#"
[global]
level = "debug"
format = "{{timestamp}} [{{level}}] {{target}} - {{message}}"
[performance]
channel_capacity = 500
"#
)
.expect("Failed to write config to temp file");
let config_content = std::fs::read_to_string(file.path()).expect("Failed to read config file");
let config: inklog::InklogConfig = config_content.parse().expect("Failed to parse config");
assert_eq!(config.global.level, "debug");
assert_eq!(config.performance.channel_capacity, 500);
assert_eq!(config.performance.worker_threads, 3);
assert!(config.validate().is_ok());
}
use inklog::LoggerManager as RecoveryLoggerManager;
use std::fs as recovery_fs;
use std::thread as recovery_thread;
use std::time::Duration as RecoveryDuration;
#[tokio::test(flavor = "multi_thread")]
async fn test_file_sink_auto_recovery() {
let test_dir = "tests/temp_recovery";
let _ = recovery_fs::create_dir_all(test_dir);
let log_file = format!("{}/test_recovery.log", test_dir);
let manager = RecoveryLoggerManager::builder()
.level("info")
.file(log_file.clone())
.build()
.await
.expect("Failed to create logger manager");
tracing::info!("Test message before failure");
recovery_thread::sleep(RecoveryDuration::from_millis(100));
let _ = recovery_fs::remove_file(&log_file);
for i in 0..10 {
tracing::info!("Test message during failure {}", i);
recovery_thread::sleep(RecoveryDuration::from_millis(50));
}
recovery_thread::sleep(RecoveryDuration::from_secs(2));
tracing::info!("Test message after recovery");
recovery_thread::sleep(RecoveryDuration::from_millis(100));
let health = manager.get_health_status();
println!("Health status: {:?}", health);
let _ = recovery_fs::remove_dir_all(test_dir);
}
#[tokio::test(flavor = "multi_thread")]
async fn test_manual_sink_recovery() {
let test_dir = "tests/temp_manual_recovery";
let _ = recovery_fs::create_dir_all(test_dir);
let log_file = format!("{}/test_manual_recovery.log", test_dir);
let manager = RecoveryLoggerManager::builder()
.level("info")
.file(log_file.clone())
.build()
.await
.expect("Failed to create logger manager");
tracing::info!("Initial test message");
recovery_thread::sleep(RecoveryDuration::from_millis(100));
let _ = recovery_fs::remove_file(&log_file);
tracing::info!("Message during failure");
recovery_thread::sleep(RecoveryDuration::from_millis(100));
let recovery_result = manager.recover_sink("file");
println!("Manual recovery result: {:?}", recovery_result);
recovery_thread::sleep(RecoveryDuration::from_millis(500));
tracing::info!("Message after manual recovery");
recovery_thread::sleep(RecoveryDuration::from_millis(100));
let _ = recovery_fs::remove_dir_all(test_dir);
assert!(recovery_result.is_ok());
}
#[tokio::test(flavor = "multi_thread")]
async fn test_bulk_recovery_for_unhealthy_sinks() {
let test_dir = "tests/temp_bulk_recovery";
let _ = recovery_fs::create_dir_all(test_dir);
let log_file = format!("{}/test_bulk_recovery.log", test_dir);
let manager = RecoveryLoggerManager::builder()
.level("info")
.file(log_file.clone())
.build()
.await
.expect("Failed to create logger manager");
tracing::info!("Initial test message");
recovery_thread::sleep(RecoveryDuration::from_millis(100));
let _ = recovery_fs::remove_file(&log_file);
for i in 0..5 {
tracing::info!("Message during failure {}", i);
recovery_thread::sleep(RecoveryDuration::from_millis(50));
}
let recovery_result = manager.trigger_recovery_for_unhealthy_sinks();
println!("Bulk recovery result: {:?}", recovery_result);
recovery_thread::sleep(RecoveryDuration::from_millis(500));
tracing::info!("Message after bulk recovery");
recovery_thread::sleep(RecoveryDuration::from_millis(100));
let _ = recovery_fs::remove_dir_all(test_dir);
assert!(recovery_result.is_ok());
}
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
use inklog::config::DatabaseDriver as BatchDatabaseDriver;
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
use inklog::log_record::LogRecord as BatchLogRecord;
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
use inklog::sink::database::DatabaseSink as BatchDatabaseSink;
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
#[allow(unused_imports)]
use inklog::sink::LogSink as BatchLogSink;
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
use inklog::DatabaseSinkConfig as BatchDatabaseSinkConfig;
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
use std::time::Duration as BatchDuration;
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
use tempfile::TempDir as BatchTempDir;
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
use tracing::Level as BatchLevel;
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
#[tokio::test(flavor = "multi_thread")]
async fn test_database_batch_write_dbnexus() {
let temp_dir = BatchTempDir::new().expect("Failed to create temp directory");
let db_path = temp_dir.path().join("logs.db");
let url = format!("sqlite://{}?mode=rwc", db_path.display());
let _ = create_logs_table(&url).await;
let config = BatchDatabaseSinkConfig {
name: "test".to_string(),
enabled: true,
driver: BatchDatabaseDriver::SQLite,
url: url.clone(),
batch_size: 5,
flush_interval_ms: 1000,
pool_size: 5,
partition: inklog::config::PartitionStrategy::default(),
table_name: "logs".to_string(),
archive_format: "json".to_string(),
parquet_config: inklog::config::ParquetConfig::default(),
};
let mock_db = inklog::integrations::infra::MockDatabaseAdapter::new();
let sink = BatchDatabaseSink::new_with_config(std::sync::Arc::new(mock_db), Some(config))
.expect("Failed to create DatabaseSink");
for i in 0..3 {
let record = BatchLogRecord::new(
BatchLevel::INFO,
"batch_test".into(),
format!("Message {}", i),
);
sink.write(&record)
.await
.expect("Failed to write log record");
}
tokio::time::sleep(BatchDuration::from_millis(1100)).await;
let record = BatchLogRecord::new(
BatchLevel::INFO,
"batch_test".into(),
"Trigger flush".into(),
);
sink.write(&record)
.await
.expect("Failed to write log record");
tokio::time::sleep(BatchDuration::from_millis(200)).await;
sink.flush().await.expect("Failed to flush batch logs");
for i in 4..9 {
let record = BatchLogRecord::new(
BatchLevel::INFO,
"batch_test".into(),
format!("Message {}", i),
);
sink.write(&record)
.await
.expect("Failed to write log record");
}
tokio::time::sleep(BatchDuration::from_millis(500)).await;
sink.flush().await.expect("Failed to flush batch logs");
}
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
#[tokio::test(flavor = "multi_thread")]
async fn test_database_timeout_flush_dbnexus() {
let temp_dir = BatchTempDir::new().expect("Failed to create temp directory");
let db_path = temp_dir.path().join("logs_timeout.db");
let url = format!("sqlite://{}?mode=rwc", db_path.display());
let _ = create_logs_table(&url).await;
let config = BatchDatabaseSinkConfig {
name: "test".to_string(),
enabled: true,
driver: BatchDatabaseDriver::SQLite,
url: url.clone(),
batch_size: 100,
flush_interval_ms: 300,
pool_size: 5,
partition: inklog::config::PartitionStrategy::default(),
table_name: "logs".to_string(),
archive_format: "json".to_string(),
parquet_config: inklog::config::ParquetConfig::default(),
};
let mock_db = inklog::integrations::infra::MockDatabaseAdapter::new();
let sink = BatchDatabaseSink::new_with_config(std::sync::Arc::new(mock_db), Some(config))
.expect("Failed to create DatabaseSink");
let record1 = BatchLogRecord::new(
BatchLevel::INFO,
"timeout_test".into(),
"First message".into(),
);
sink.write(&record1)
.await
.expect("Failed to write first log record");
tokio::time::sleep(BatchDuration::from_millis(500)).await;
let record2 = BatchLogRecord::new(
BatchLevel::INFO,
"timeout_test".into(),
"Second message".into(),
);
sink.write(&record2)
.await
.expect("Failed to write second log record");
tokio::time::sleep(BatchDuration::from_millis(500)).await;
sink.flush().await.expect("Failed to flush timeout logs");
}
use inklog::InklogConfig as ConfigInklogConfig;
use serial_test::serial as config_serial;
fn clear_all_inklog_env_vars() {
for (key, _) in std::env::vars() {
if key.starts_with("INKLOG_") {
std::env::remove_var(&key);
}
}
}
#[test]
#[config_serial]
fn test_config_from_env_overrides() {
clear_all_inklog_env_vars();
std::env::set_var("INKLOG_GLOBAL_LEVEL", "debug");
std::env::set_var("INKLOG_FILE_SINK_ENABLED", "true");
std::env::set_var("INKLOG_FILE_SINK_PATH", "/tmp/test_logs/app.log");
std::env::set_var("INKLOG_FILE_SINK_MAX_SIZE", "50MB");
std::env::set_var("INKLOG_FILE_SINK_COMPRESS", "true");
let config = ConfigInklogConfig::load_with_env_overrides().unwrap();
assert_eq!(config.global.level, "debug");
assert!(config.file_sink.is_some());
let file = config.file_sink.unwrap();
assert!(file.enabled);
assert_eq!(file.max_size, "50MB");
assert!(file.compress);
}
#[test]
#[config_serial]
fn test_config_env_override_http_server() {
clear_all_inklog_env_vars();
std::env::set_var("INKLOG_HTTP_SERVER_ENABLED", "true");
std::env::set_var("INKLOG_HTTP_SERVER_HOST", "127.0.0.1");
std::env::set_var("INKLOG_HTTP_SERVER_PORT", "9090");
std::env::set_var("INKLOG_HTTP_SERVER_METRICS_PATH", "/prometheus");
std::env::set_var("INKLOG_HTTP_SERVER_HEALTH_PATH", "/status");
let config = ConfigInklogConfig::load_with_env_overrides().unwrap();
assert!(config.http_server.is_some());
let http = config.http_server.unwrap();
assert!(http.enabled);
assert_eq!(http.host, "127.0.0.1");
assert_eq!(http.port, 9090);
assert_eq!(http.metrics_path, "/prometheus");
assert_eq!(http.health_path, "/status");
}
#[test]
#[config_serial]
fn test_config_env_override_performance() {
clear_all_inklog_env_vars();
std::env::set_var("INKLOG_PERFORMANCE_WORKER_THREADS", "8");
std::env::set_var("INKLOG_PERFORMANCE_CHANNEL_CAPACITY", "20000");
let config = ConfigInklogConfig::load_with_env_overrides().unwrap();
assert_eq!(config.performance.worker_threads, 8);
assert_eq!(config.performance.channel_capacity, 20000);
}
use inklog::config::{HttpErrorMode, HttpServerConfig};
use inklog::InklogConfig as HttpInklogConfig;
use serial_test::serial as http_serial;
fn clear_inklog_env() {
for (key, _) in std::env::vars() {
if key.starts_with("INKLOG_") {
std::env::remove_var(&key);
}
}
}
#[tokio::test]
#[http_serial]
async fn test_http_server_startup_with_default_config() {
clear_inklog_env();
let port = 18080
+ std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs() as u16
% 10000;
let config = HttpServerConfig {
enabled: true,
host: "127.0.0.1".to_string(),
port,
metrics_path: "/metrics".to_string(),
health_path: "/health".to_string(),
error_mode: HttpErrorMode::Strict,
auth: None,
ip_whitelist: None,
};
let inklog_config = HttpInklogConfig {
http_server: Some(config),
..Default::default()
};
assert!(inklog_config.http_server.is_some());
let http = inklog_config.http_server.unwrap();
assert!(http.enabled);
assert_eq!(http.port, port);
}
#[tokio::test]
#[http_serial]
async fn test_http_server_error_mode_panic() {
clear_inklog_env();
let config = HttpServerConfig {
enabled: true,
host: "127.0.0.1".to_string(),
port: 18081,
metrics_path: "/metrics".to_string(),
health_path: "/health".to_string(),
error_mode: HttpErrorMode::Strict,
auth: None,
ip_whitelist: None,
};
match config.error_mode {
HttpErrorMode::Strict => {}
_ => panic!("Expected Strict mode"),
}
}
#[tokio::test]
#[http_serial]
async fn test_http_server_error_mode_warn() {
clear_inklog_env();
let config = HttpServerConfig {
enabled: true,
host: "127.0.0.1".to_string(),
port: 18082,
metrics_path: "/metrics".to_string(),
health_path: "/health".to_string(),
error_mode: HttpErrorMode::Warn,
auth: None,
ip_whitelist: None,
};
match config.error_mode {
HttpErrorMode::Warn => {}
_ => panic!("Expected Warn mode"),
}
}
#[tokio::test]
#[http_serial]
async fn test_http_server_error_mode_strict() {
clear_inklog_env();
let config = HttpServerConfig {
enabled: true,
host: "127.0.0.1".to_string(),
port: 18083,
metrics_path: "/metrics".to_string(),
health_path: "/health".to_string(),
error_mode: HttpErrorMode::Strict,
auth: None,
ip_whitelist: None,
};
match config.error_mode {
HttpErrorMode::Strict => {}
_ => panic!("Expected Strict mode"),
}
}
#[http_serial]
#[tokio::test]
async fn test_http_server_with_logger_manager() {
clear_inklog_env();
std::env::set_var("INKLOG_HTTP_SERVER_ENABLED", "true");
std::env::set_var("INKLOG_HTTP_SERVER_HOST", "127.0.0.1");
std::env::set_var("INKLOG_HTTP_SERVER_PORT", "18084");
std::env::set_var("INKLOG_HTTP_SERVER_ERROR_MODE", "warn");
let config = HttpInklogConfig::load_with_env_overrides().unwrap();
assert!(config.http_server.is_some());
let http = config.http_server.unwrap();
assert!(http.enabled);
assert_eq!(http.host, "127.0.0.1");
assert_eq!(http.port, 18084);
match http.error_mode {
HttpErrorMode::Warn => {}
_ => panic!("Expected Warn mode from env"),
}
std::env::remove_var("INKLOG_HTTP_SERVER_ENABLED");
std::env::remove_var("INKLOG_HTTP_SERVER_HOST");
std::env::remove_var("INKLOG_HTTP_SERVER_PORT");
std::env::remove_var("INKLOG_HTTP_SERVER_ERROR_MODE");
}
#[http_serial]
#[tokio::test]
async fn test_http_metrics_path_configuration() {
clear_inklog_env();
std::env::set_var("INKLOG_HTTP_SERVER_ENABLED", "true");
std::env::set_var("INKLOG_HTTP_SERVER_METRICS_PATH", "/prometheus/metrics");
std::env::set_var("INKLOG_HTTP_SERVER_HEALTH_PATH", "/status");
let config = HttpInklogConfig::load_with_env_overrides().unwrap();
let http = config
.http_server
.expect("http_server should be Some after setting INKLOG_HTTP_SERVER_ENABLED");
assert_eq!(http.metrics_path, "/prometheus/metrics");
assert_eq!(http.health_path, "/status");
}
#[tokio::test]
#[http_serial]
async fn test_http_server_disabled_by_default() {
clear_inklog_env();
let config = HttpInklogConfig::load_sync().unwrap();
assert!(
config.http_server.is_none(),
"INKLOG_HTTP_SERVER_ENABLED should not be set"
);
}
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
use arrow_array::RecordBatchReader;
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
use arrow_schema::DataType;
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
use bytes::Bytes;
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
use inklog::sink::database::convert_logs_to_parquet;
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
use std::time::Instant;
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
fn create_test_logs(count: usize) -> Vec<inklog::log_record::LogRecord> {
(0..count)
.map(|i| inklog::log_record::LogRecord {
timestamp: chrono::Utc::now(),
level: match i % 5 {
0 => "trace".to_string(),
1 => "debug".to_string(),
2 => "info".to_string(),
3 => "warn".to_string(),
_ => "error".to_string(),
},
target: format!("test_module::function_{}", i % 10),
message: format!("Test log message number {}", i),
fields: std::collections::HashMap::new(),
file: Some(format!("src/test_{}.rs", i % 5)),
line: Some((i % 100) as u32),
thread_id: format!("thread-{}", i % 4),
})
.collect()
}
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
const EXPECTED_FIELD_NAMES: &[&str] = &[
"id",
"timestamp",
"level",
"target",
"message",
"fields",
"file",
"line",
"thread_id",
];
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
const EXPECTED_FIELD_TYPES: &[DataType] = &[
DataType::Int64, DataType::Date64, DataType::Utf8, DataType::Utf8, DataType::Utf8, DataType::Utf8, DataType::Utf8, DataType::Int32, DataType::Utf8, ];
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
fn verify_parquet_schema(data: &[u8]) -> Result<(), Box<dyn std::error::Error>> {
let bytes = Bytes::copy_from_slice(data);
let reader = ParquetRecordBatchReaderBuilder::try_new(bytes)?.build()?;
let schema = reader.schema();
let fields = schema.fields();
assert_eq!(fields.len(), 9, "Schema should have 9 fields");
for (i, (name, dtype)) in EXPECTED_FIELD_NAMES
.iter()
.zip(EXPECTED_FIELD_TYPES.iter())
.enumerate()
{
assert_eq!(fields[i].name(), *name);
assert_eq!(fields[i].data_type(), dtype);
}
Ok(())
}
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
fn verify_parquet_data(data: &[u8]) -> Result<(), Box<dyn std::error::Error>> {
let bytes = Bytes::copy_from_slice(data);
let reader = ParquetRecordBatchReaderBuilder::try_new(bytes)?.build()?;
let mut total_rows = 0;
for batch in reader {
let batch = batch?;
assert!(batch.num_rows() > 0, "Batch should have rows");
total_rows += batch.num_rows();
}
assert!(total_rows > 0, "Parquet file should contain data");
Ok(())
}
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
fn verify_parquet_file(data: &[u8]) -> Result<(), Box<dyn std::error::Error>> {
verify_parquet_schema(data)?;
verify_parquet_data(data)?;
Ok(())
}
#[test]
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
fn test_parquet_basic_conversion() {
let logs = create_test_logs(100);
let result = convert_logs_to_parquet(&logs, &Default::default());
assert!(
result.is_ok(),
"Parquet conversion should succeed: {}",
result.unwrap_err()
);
let parquet_data = result.expect("Parquet conversion should succeed");
assert!(!parquet_data.is_empty(), "Parquet data should not be empty");
verify_parquet_file(&parquet_data).expect("Parquet file should be valid");
}
#[test]
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
fn test_parquet_small_dataset() {
let logs = create_test_logs(1_000);
let start = Instant::now();
let result = convert_logs_to_parquet(&logs, &Default::default());
let duration = start.elapsed();
let parquet_data = result.expect("Parquet conversion should succeed for 1K records");
println!("1K records conversion time: {:?}", duration);
println!("1K records Parquet size: {} bytes", parquet_data.len());
let estimated_original_size = logs.len() * 200;
let compression_ratio = estimated_original_size as f64 / parquet_data.len() as f64;
println!("Estimated compression ratio: {:.2}x", compression_ratio);
assert!(
compression_ratio > 1.5,
"Compression ratio should be > 1.5x, got {:.2}x",
compression_ratio
);
verify_parquet_file(&parquet_data).expect("Parquet file should be valid");
}
#[test]
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
fn test_parquet_medium_dataset() {
let logs = create_test_logs(10_000);
let start = Instant::now();
let result = convert_logs_to_parquet(&logs, &Default::default());
let duration = start.elapsed();
let parquet_data = result.expect("Parquet conversion should succeed for 10K records");
println!("10K records conversion time: {:?}", duration);
println!("10K records Parquet size: {} bytes", parquet_data.len());
assert!(
duration.as_secs() < 5,
"10K records conversion should complete in < 5 seconds, took {:?}",
duration
);
verify_parquet_file(&parquet_data).expect("Parquet file should be valid");
}
#[test]
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
fn test_parquet_large_dataset() {
let logs = create_test_logs(100_000);
let start = Instant::now();
let result = convert_logs_to_parquet(&logs, &Default::default());
let duration = start.elapsed();
let parquet_data = result.expect("Parquet conversion should succeed for 100K records");
println!("100K records conversion time: {:?}", duration);
println!("100K records Parquet size: {} bytes", parquet_data.len());
assert!(
duration.as_secs() < 30,
"100K records conversion should complete in < 30 seconds, took {:?}",
duration
);
verify_parquet_file(&parquet_data).expect("Parquet file should be valid");
}
#[test]
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
fn test_parquet_compression_ratio() {
let logs = create_test_logs(10_000);
let result = convert_logs_to_parquet(&logs, &Default::default())
.expect("Parquet conversion should succeed");
let json_data = serde_json::to_vec(&logs).expect("JSON serialization should succeed");
let original_size = json_data.len();
let compressed_size = result.len();
let compression_ratio = original_size as f64 / compressed_size as f64;
println!("Original JSON size: {} bytes", original_size);
println!("Compressed Parquet size: {} bytes", compressed_size);
println!("Actual compression ratio: {:.2}x", compression_ratio);
assert!(
compression_ratio > 2.0,
"Compression ratio should be > 2.0x, got {:.2}x",
compression_ratio
);
}
#[test]
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
fn test_parquet_empty_dataset() {
let logs: Vec<inklog::log_record::LogRecord> = vec![];
let result = convert_logs_to_parquet(&logs, &Default::default());
let parquet_data = result.expect("Parquet conversion should succeed for empty dataset");
assert!(
!parquet_data.is_empty(),
"Parquet file should have metadata even for empty data"
);
}
#[test]
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
fn test_parquet_schema_compatibility() {
let logs = create_test_logs(100);
let result = convert_logs_to_parquet(&logs, &Default::default())
.expect("Parquet conversion should succeed");
verify_parquet_schema(&result).expect("Schema verification should pass");
}
use inklog::LoggerManager as StabilityLoggerManager;
use std::thread as stability_thread;
use std::time::{Duration as StabilityDuration, Instant as StabilityInstant};
use tracing::{error as stability_error, info as stability_info};
#[tokio::test(flavor = "multi_thread")]
#[ignore = "manual"] async fn test_long_running_stability() {
let logger = StabilityLoggerManager::builder()
.level("debug")
.build()
.await
.expect("Failed to create LoggerManager");
let duration = StabilityDuration::from_secs(5); let start = StabilityInstant::now();
let handles: Vec<_> = (0..4)
.map(|i| {
stability_thread::spawn(move || {
let mut count = 0;
while start.elapsed() < duration {
stability_info!(target: "stability", "Thread {} log {}", i, count);
if count % 100 == 0 {
stability_error!(target: "stability", "Thread {} error {}", i, count);
}
count += 1;
stability_thread::sleep(StabilityDuration::from_millis(1));
}
})
})
.collect();
for h in handles {
h.join().expect("Thread join failed");
}
stability_thread::sleep(StabilityDuration::from_millis(500));
let status = logger.get_health_status();
println!(
"Stability test passed. Status: {:?}, Metrics: {:?}",
status.overall_status, status.metrics
);
}
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
use inklog::config::DatabaseDriver as VerifyDatabaseDriver;
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
use inklog::sink::database::DatabaseSink as VerifyDatabaseSink;
use inklog::sink::file::FileSink as VerifyFileSink;
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
use inklog::{
log_record::LogRecord as VerifyLogRecord, DatabaseSinkConfig as VerifyDatabaseSinkConfig,
FileSinkConfig as VerifyFileSinkConfig,
};
#[cfg(not(any(feature = "sqlite", feature = "postgres", feature = "mysql")))]
use inklog::{log_record::LogRecord as VerifyLogRecord, FileSinkConfig as VerifyFileSinkConfig};
use std::fs::File as VerifyFile;
use std::io::Read as VerifyRead;
use std::path::PathBuf;
use std::time::Duration as VerifyDuration;
use tempfile::TempDir as VerifyTempDir;
use tracing::Level as VerifyLevel;
fn find_file_with_extension(dir: &VerifyTempDir, extension: &str) -> Option<PathBuf> {
let paths: Vec<_> = std::fs::read_dir(dir.path())
.expect("Failed to read temp directory")
.filter_map(|entry| entry.ok())
.map(|e| e.path())
.collect();
paths
.into_iter()
.find(|p| p.extension().is_some_and(|ext| ext == extension))
}
fn verify_zstd_compression(file_path: &PathBuf) {
let mut file = VerifyFile::open(file_path).expect("Failed to open compressed file");
let mut magic = [0u8; 4];
file.read_exact(&mut magic)
.expect("Failed to read file magic bytes");
assert_eq!(magic, [0x28, 0xB5, 0x2F, 0xFD]);
}
fn verify_encrypted_file(file_path: &PathBuf) {
let metadata = std::fs::metadata(file_path).expect("Failed to get file metadata");
assert!(
metadata.len() > 12,
"Encrypted file should have nonce (12 bytes) + ciphertext"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn verify_file_sink_compression() {
let temp_dir = VerifyTempDir::new().expect("Failed to create temp directory");
let log_path = temp_dir.path().join("test.log");
let config = VerifyFileSinkConfig {
enabled: true,
path: log_path.clone(),
max_size: "10".into(),
compress: true,
encrypt: false,
..Default::default()
};
let sink = VerifyFileSink::new(config).expect("Failed to create FileSink");
let record = VerifyLogRecord::new(
VerifyLevel::INFO,
"test".into(),
"A long message to trigger rotation".into(),
);
sink.write(&record)
.await
.expect("Failed to write log record");
for _ in 0..5 {
sink.write(&record)
.await
.expect("Failed to write log record during rotation");
}
std::thread::sleep(VerifyDuration::from_millis(1000));
let zst_path = find_file_with_extension(&temp_dir, "zst").expect("No compressed file found");
verify_zstd_compression(&zst_path);
}
#[tokio::test(flavor = "multi_thread")]
async fn verify_file_sink_encryption() {
let temp_dir = VerifyTempDir::new().expect("Failed to create temp directory");
let log_path = temp_dir.path().join("enc.log");
std::env::set_var("LOG_KEY", "YWJjZGVmZ2hpamtsbW5vcHFyc3R1dnd4eXoxMjM0NTY=");
let config = VerifyFileSinkConfig {
enabled: true,
path: log_path.clone(),
max_size: "100".into(),
compress: false,
encrypt: true,
encryption_key_env: Some("LOG_KEY".into()),
..Default::default()
};
let sink = VerifyFileSink::new(config).expect("Failed to create FileSink");
let record = VerifyLogRecord::new(VerifyLevel::INFO, "test".into(), "Secret message".into());
sink.write(&record)
.await
.expect("Failed to write log record");
for _ in 0..5 {
sink.write(&record)
.await
.expect("Failed to write log record during rotation");
}
sink.flush().await.expect("Failed to flush");
std::thread::sleep(VerifyDuration::from_millis(1000));
let enc_path = find_file_with_extension(&temp_dir, "enc").expect("No encrypted file found");
verify_encrypted_file(&enc_path);
}
#[tokio::test(flavor = "multi_thread")]
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
async fn verify_database_sink_sqlite() {
let temp_dir = VerifyTempDir::new().expect("Failed to create temp directory");
let db_path = temp_dir.path().join("logs.db");
let url = format!("sqlite://{}?mode=rwc", db_path.display());
let _ = create_logs_table(&url).await;
let config = VerifyDatabaseSinkConfig {
enabled: true,
driver: VerifyDatabaseDriver::SQLite,
url: url.clone(),
batch_size: 1,
flush_interval_ms: 100,
..Default::default()
};
let mock_db = inklog::integrations::infra::MockDatabaseAdapter::new();
let mock_db_arc = std::sync::Arc::new(mock_db);
let sink = VerifyDatabaseSink::new_with_config(mock_db_arc.clone(), Some(config))
.expect("Failed to create DatabaseSink");
let record = VerifyLogRecord::new(VerifyLevel::INFO, "db_test".into(), "message to db".into());
sink.write(&record)
.await
.expect("Failed to write log record to database");
tokio::time::sleep(VerifyDuration::from_millis(500)).await;
sink.flush().await.expect("Failed to flush database sink");
let mock_ref = mock_db_arc.as_ref() as &inklog::integrations::infra::MockDatabaseAdapter;
assert_eq!(mock_ref.record_count(), 1);
{
use inklog::sink::entity::{
sea_orm::{Database, EntityTrait},
Entity,
};
let db = Database::connect(&url)
.await
.expect("Failed to connect to database");
let logs = Entity::find().all(&db).await.expect("Failed to query logs");
let _ = logs;
}
}
#[cfg(any(feature = "sqlite", feature = "postgres", feature = "mysql"))]
async fn create_logs_table(url: &str) -> Result<(), String> {
let pool = dbnexus::DbPool::new(url).await.map_err(|e| e.to_string())?;
let session = pool.get_session("admin").await.map_err(|e| e.to_string())?;
use inklog::sink::entity::sea_orm::{ConnectionTrait, Schema};
let conn = session.connection().map_err(|e| e.to_string())?;
let schema = Schema::new(conn.get_database_backend());
conn.execute(
schema
.create_table_from_entity(inklog::sink::entity::Entity)
.if_not_exists(),
)
.await
.map_err(|e: inklog::sink::entity::sea_orm::DbErr| e.to_string())?;
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn test_log_crate_native_support() {
let _logger = LoggerManager::builder().level("debug").build().await;
log::info!("This is a log::info message");
log::warn!("This is a log::warn message");
log::error!("This is a log::error message");
log::debug!("This is a log::debug message");
std::thread::sleep(Duration::from_millis(200));
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn test_tracing_and_log_coexist() {
let _logger = LoggerManager::builder().level("debug").build().await;
log::info!("log::info message");
tracing::info!("tracing::info message");
log::error!("log::error message");
tracing::error!("tracing::error message");
std::thread::sleep(Duration::from_millis(200));
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn test_log_level_filtering() {
let _logger = LoggerManager::builder().level("warn").build().await;
log::debug!("This debug message should not appear");
log::info!("This info message should not appear");
log::warn!("This warn message should appear");
log::error!("This error message should appear");
std::thread::sleep(Duration::from_millis(100));
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn test_log_all_levels() {
let _logger = LoggerManager::builder().level("trace").build().await;
log::trace!("Trace message from log crate");
log::debug!("Debug message from log crate");
log::info!("Info message from log crate");
log::warn!("Warn message from log crate");
log::error!("Error message from log crate");
std::thread::sleep(Duration::from_millis(100));
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn test_log_to_file() {
let temp_dir = tempfile::tempdir().unwrap();
let log_file = temp_dir.path().join("test.log");
let logger = match LoggerManager::builder()
.level("info")
.file(&log_file)
.build()
.await
{
Ok(l) => l,
Err(_) => {
println!("LoggerManager init failed, skipping test_log_to_file");
return;
}
};
log::info!("PROBE_FILE_LOG");
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
let probe_ok = log_file.exists()
&& std::fs::read_to_string(&log_file)
.map(|c| c.contains("PROBE_FILE_LOG"))
.unwrap_or(false);
if !probe_ok {
println!("Global logger not effective, skipping test_log_to_file");
drop(logger);
return;
}
log::info!("This should go to file");
log::warn!("This warning should also be in file");
tokio::time::sleep(Duration::from_millis(500)).await;
assert!(log_file.exists(), "Log file should exist");
let contents = std::fs::read_to_string(&log_file).unwrap_or_default();
if contents.is_empty() {
println!("Warning: Log file is empty, logger may not have initialized properly");
}
let _ = logger.shutdown();
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn test_concurrent_file_writes() {
use inklog::{FileSinkConfig, InklogConfig};
use std::sync::{Arc, Barrier};
use std::thread;
use tempfile::TempDir;
let temp_dir = TempDir::new().unwrap();
let log_path = temp_dir.path().join("concurrent_test.log");
let config = InklogConfig {
file_sink: Some(FileSinkConfig {
enabled: true,
path: log_path.clone(),
max_size: "100MB".into(),
batch_size: 100,
flush_interval_ms: 100,
..Default::default()
}),
performance: inklog::config::PerformanceConfig {
worker_threads: 4,
channel_capacity: 10000,
..Default::default()
},
..Default::default()
};
let logger = match LoggerManager::with_config(config).await {
Ok(l) => l,
Err(_) => {
println!("LoggerManager init failed, skipping test_concurrent_file_writes");
return;
}
};
log::info!(target: "concurrent_test", "PROBE_MESSAGE");
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
let probe_ok = log_path.exists()
&& std::fs::read_to_string(&log_path)
.map(|c| c.contains("PROBE_MESSAGE"))
.unwrap_or(false);
if !probe_ok {
println!("Global logger not effective (already set by other test), skipping test_concurrent_file_writes");
drop(logger);
return;
}
let num_threads = 4;
let messages_per_thread = 100;
let barrier = Arc::new(Barrier::new(num_threads));
let handles: Vec<_> = (0..num_threads)
.map(|thread_id| {
let barrier = Arc::clone(&barrier);
thread::spawn(move || {
barrier.wait();
for i in 0..messages_per_thread {
log::info!(target: "concurrent_test", "Thread {} - Message {}", thread_id, i);
}
})
})
.collect();
for handle in handles {
handle.join().unwrap();
}
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
assert!(log_path.exists());
let metadata = std::fs::metadata(&log_path).unwrap();
assert!(
metadata.len() > 1000,
"Expected file > 1000 bytes, got {}",
metadata.len()
);
let _ = logger.shutdown();
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn test_memory_stability() {
use inklog::{FileSinkConfig, InklogConfig};
use tempfile::TempDir;
let temp_dir = TempDir::new().unwrap();
let log_path = temp_dir.path().join("memory_test.log");
let config = InklogConfig {
file_sink: Some(FileSinkConfig {
enabled: true,
path: log_path.clone(),
max_size: "100MB".into(),
batch_size: 100,
flush_interval_ms: 100,
..Default::default()
}),
..Default::default()
};
let logger = match LoggerManager::with_config(config).await {
Ok(l) => l,
Err(_) => {
println!("LoggerManager init failed, skipping test_memory_stability");
return;
}
};
log::info!(target: "memory_test", "PROBE_MESSAGE");
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
let probe_ok = log_path.exists()
&& std::fs::read_to_string(&log_path)
.map(|c| c.contains("PROBE_MESSAGE"))
.unwrap_or(false);
if !probe_ok {
println!("Global logger not effective (already set by other test), skipping test_memory_stability");
drop(logger);
return;
}
for i in 0..1000 {
log::info!(target: "memory_test", "Memory test message {}", i);
}
tokio::time::sleep(std::time::Duration::from_secs(3)).await;
assert!(log_path.exists());
let metadata = std::fs::metadata(&log_path).unwrap();
assert!(
metadata.len() > 5000,
"Expected file > 5000 bytes, got {}",
metadata.len()
);
let _ = logger.shutdown();
}