use crate::support::io::sink::circuit_breaker::CircuitBreaker;
use crate::support::io::sink::rotation::{RotationStrategy, SizeBasedRotation, TimeBasedRotation};
use crate::support::io::sink::LogSink;
use crate::DataMasker;
use crate::FileSinkConfig;
use crate::InklogError;
use crate::LogRecord;
use aes_gcm::aead::Aead;
use aes_gcm::KeyInit;
use async_trait::async_trait;
use bytes::BytesMut;
use chrono::{DateTime, Datelike, Utc};
use parking_lot::RwLock;
use std::fs::{self, File, OpenOptions};
use std::io::Write;
use std::path::Path;
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::thread;
use std::time::{Duration as StdDuration, Instant};
use tracing::{debug, error, info, warn};
pub use super::circuit_breaker::{CircuitBreakerConfig, CircuitState};
struct FileSinkInner {
current_file: Option<File>,
current_size: u64,
last_rotation: Instant,
next_rotation_time: Option<DateTime<Utc>>,
last_rotation_date: Option<i32>,
sequence: u32,
batch_buffer: Vec<LogRecord>,
last_flush_time: Instant,
circuit_breaker: CircuitBreaker,
fallback_sink: Option<Arc<dyn LogSink + Send + Sync>>,
rotation_timer: Option<Arc<parking_lot::Mutex<Instant>>>,
timer_handle: Option<thread::JoinHandle<()>>,
cleanup_timer_handle: Option<thread::JoinHandle<()>>,
rotation_strategy: Box<dyn RotationStrategy>,
}
pub struct FileSink {
config: FileSinkConfig,
rotation_interval: StdDuration,
last_cleanup_time: Arc<parking_lot::Mutex<Option<Instant>>>,
shutdown_flag: Arc<AtomicBool>,
masker: DataMasker,
inner: RwLock<FileSinkInner>,
}
impl FileSink {
pub fn new(config: FileSinkConfig) -> Result<Self, InklogError> {
let rotation_interval = match config.rotation_time.as_str() {
"hourly" => StdDuration::from_secs(3600),
"daily" => StdDuration::from_secs(86400),
"weekly" => StdDuration::from_secs(604800),
"monthly" => StdDuration::from_secs(2592000),
_ => StdDuration::from_secs(86400),
};
let rotation_timer = Arc::new(parking_lot::Mutex::new(Instant::now()));
let last_rotation = Instant::now();
let rotation_strategy: Box<dyn RotationStrategy> = {
let max_size = Self::parse_size(&config.max_size).unwrap_or(100 * 1024 * 1024);
let size_strategy = SizeBasedRotation::new(max_size);
let time_strategy = TimeBasedRotation::from_interval_string(&config.rotation_time)
.unwrap_or_else(|_| {
TimeBasedRotation::from_interval_string("daily")
.expect("hardcoded 'daily' interval is valid")
});
Box::new(crate::support::io::sink::rotation::CompositeRotation::new(
vec![Box::new(size_strategy), Box::new(time_strategy)],
))
};
let inner = FileSinkInner {
current_file: None,
current_size: 0,
last_rotation,
next_rotation_time: None,
last_rotation_date: None,
sequence: 0,
fallback_sink: None,
circuit_breaker: CircuitBreaker::new(5, StdDuration::from_secs(30), 3),
batch_buffer: Vec::with_capacity(config.batch_size),
last_flush_time: Instant::now(),
timer_handle: None,
rotation_timer: Some(rotation_timer.clone()),
cleanup_timer_handle: None,
rotation_strategy,
};
let sink = Self {
config: config.clone(),
rotation_interval,
last_cleanup_time: Arc::new(parking_lot::Mutex::new(None)),
shutdown_flag: Arc::new(AtomicBool::new(false)),
masker: DataMasker::new(),
inner: RwLock::new(inner),
};
{
let mut inner = sink.inner.write();
sink.update_next_rotation_time_inner(&mut inner);
}
{
let mut inner = sink.inner.write();
if let Err(e) = sink.open_file_inner(&mut inner) {
error!("Failed to open log file: {}", e);
return Err(e);
}
}
sink.start_rotation_timer();
sink.start_cleanup_timer();
Ok(sink)
}
pub fn parse_size(size_str: &str) -> Option<u64> {
let size_str = size_str.trim();
if size_str.is_empty() {
return None;
}
if size_str.ends_with("TB") {
size_str
.trim_end_matches("TB")
.parse::<u64>()
.ok()
.map(|s| s * 1024 * 1024 * 1024 * 1024)
} else if size_str.ends_with("GB") {
size_str
.trim_end_matches("GB")
.parse::<u64>()
.ok()
.map(|s| s * 1024 * 1024 * 1024)
} else if size_str.ends_with("MB") {
size_str
.trim_end_matches("MB")
.parse::<u64>()
.ok()
.map(|s| s * 1024 * 1024)
} else if size_str.ends_with("KB") {
size_str
.trim_end_matches("KB")
.parse::<u64>()
.ok()
.map(|s| s * 1024)
} else {
size_str.parse::<u64>().ok()
}
}
fn get_encryption_key(&self) -> Result<BytesMut, InklogError> {
let default_key = "LOG_ENCRYPTION_KEY".to_string();
let key_str = self
.config
.encryption_key_env
.as_ref()
.unwrap_or(&default_key);
let key = std::env::var(key_str).map_err(|_| {
InklogError::EncryptionError(format!(
"Encryption key not found in environment variable: {}",
key_str
))
})?;
if key.len() < 16 {
return Err(InklogError::EncryptionError(
"Encryption key must be at least 16 characters".to_string(),
));
}
let decoded = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, &key)
.map_err(|_| {
InklogError::EncryptionError(
"Invalid base64 encoding in encryption key".to_string(),
)
})?;
if decoded.len() != 32 {
return Err(InklogError::EncryptionError(
"Encryption key must be 32 bytes (256 bits)".to_string(),
));
}
Self::validate_key_entropy(&decoded)?;
let key_bytes = BytesMut::from(&decoded[..]);
Ok(key_bytes)
}
fn validate_key_entropy(key: &[u8]) -> Result<(), InklogError> {
if key.is_empty() {
return Err(InklogError::EncryptionError(
"Encryption key cannot be empty".to_string(),
));
}
let mut freq = [0u32; 256];
for &b in key {
freq[b as usize] += 1;
}
let len = key.len() as f64;
let entropy: f64 = freq
.iter()
.filter(|&&count| count > 0)
.map(|&count| {
let p = count as f64 / len;
-p * p.log2()
})
.sum();
const MIN_ENTROPY_THRESHOLD: f64 = 4.0;
if entropy < MIN_ENTROPY_THRESHOLD {
return Err(InklogError::EncryptionError(format!(
"Encryption key has insufficient entropy ({} < {}). \
Please use a cryptographically random key.",
entropy, MIN_ENTROPY_THRESHOLD
)));
}
Ok(())
}
fn open_file_inner(&self, inner: &mut FileSinkInner) -> Result<(), InklogError> {
if let Some(parent) = self.config.path.parent() {
if let Err(e) = fs::create_dir_all(parent) {
error!("Failed to create log directory {}: {}", parent.display(), e);
return Err(InklogError::IoError(e));
}
}
match OpenOptions::new()
.create(true)
.append(true)
.open(&self.config.path)
{
Ok(file) => {
inner.current_file = Some(file);
inner.current_size = self.config.path.metadata().map(|m| m.len()).unwrap_or(0);
debug!(
"Opened log file: {} (size: {} bytes)",
self.config.path.display(),
inner.current_size
);
Ok(())
}
Err(e) => {
error!("Failed to open log file: {}", e);
Err(InklogError::IoError(e))
}
}
}
fn start_cleanup_timer(&self) {
let interval_minutes = self.config.cleanup_interval_minutes;
let cleanup_interval = StdDuration::from_secs(interval_minutes * 60);
let shutdown_flag = self.shutdown_flag.clone();
let config = self.config.clone();
let path = self.config.path.clone();
let last_cleanup_time = self.last_cleanup_time.clone();
let handle = thread::spawn(move || {
let check_interval = StdDuration::from_secs(60);
loop {
if shutdown_flag.load(Ordering::Relaxed) {
break;
}
let mut elapsed = StdDuration::ZERO;
const POLL_INTERVAL: StdDuration = StdDuration::from_millis(100);
while elapsed < check_interval {
if shutdown_flag.load(Ordering::Relaxed) {
break;
}
let step = std::cmp::min(POLL_INTERVAL, check_interval - elapsed);
thread::sleep(step);
elapsed += step;
}
if shutdown_flag.load(Ordering::Relaxed) {
break;
}
let mut last_cleanup = last_cleanup_time.lock();
let now = Instant::now();
if last_cleanup.is_none_or(|t| now.duration_since(t) >= cleanup_interval) {
if let Err(e) = Self::perform_cleanup(&config, &path) {
error!("Cleanup failed: {}", e);
} else {
*last_cleanup = Some(now);
}
}
}
});
self.inner.write().cleanup_timer_handle = Some(handle);
}
#[allow(dead_code)]
fn perform_cleanup(config: &FileSinkConfig, log_path: &Path) -> Result<(), InklogError> {
if let Some(parent) = log_path.parent() {
let entries: Result<Vec<_>, _> = fs::read_dir(parent)?.collect();
if let Ok(entries) = entries {
let cutoff_date = Utc::now()
.checked_sub_signed(chrono::Duration::days(config.retention_days as i64))
.unwrap_or_else(Utc::now);
let mut expired_count = 0;
let mut total_size = 0u64;
for entry in &entries {
total_size += entry.path().metadata()?.len();
if let Ok(modified) = entry.path().metadata().and_then(|m| m.modified()) {
let modified_utc: DateTime<Utc> = modified.into();
if modified_utc < cutoff_date {
expired_count += 1;
}
}
}
if let Some(max_total_size_bytes) = Self::parse_size(&config.max_total_size) {
if total_size > max_total_size_bytes {
let excess_size = total_size.saturating_sub(max_total_size_bytes);
let mut deleted_size: u64 = 0;
for entry in entries {
if deleted_size >= excess_size {
break;
}
if let Ok(metadata) = entry.path().metadata() {
deleted_size += metadata.len();
}
if let Err(e) = fs::remove_file(entry.path()) {
error!("Failed to remove {}: {}", entry.path().display(), e);
}
}
} else if expired_count > 0 {
let to_delete =
(entries.len() as i32 - config.keep_files as i32).max(0) as usize;
for entry in entries.into_iter().take(to_delete) {
let _ = fs::remove_file(entry.path());
}
}
}
}
}
Ok(())
}
pub fn get_disk_space_info(&self) -> Result<(u64, u64), InklogError> {
if let Some(parent) = self.config.path.parent() {
if let Ok(_metadata) = fs::metadata(parent) {
if let Ok(stat) = nix::sys::statfs::statfs(parent) {
let total_blocks = stat.blocks();
let available_blocks = stat.blocks_available();
let block_size = stat.block_size() as u64;
let total_bytes = total_blocks * block_size;
let available_bytes = available_blocks * block_size;
return Ok((total_bytes, available_bytes));
}
}
}
Err(InklogError::IoError(std::io::Error::new(
std::io::ErrorKind::NotFound,
"Unable to get disk space info",
)))
}
fn check_disk_space(&self) -> Result<bool, InklogError> {
let (_total, available) = self.get_disk_space_info()?;
let reserved = (50 * 1024 * 1024u64).max(available / 10);
Ok(available > reserved)
}
fn calculate_next_rotation_time(rotation_time: &str) -> Option<DateTime<Utc>> {
let now = Utc::now();
match rotation_time {
"hourly" => Some(now + chrono::Duration::hours(1)),
"daily" => {
let next_naive = now.date_naive().and_hms_opt(0, 0, 0)? + chrono::Duration::days(1);
Some(next_naive.and_utc())
}
"weekly" => {
let next_naive =
now.date_naive().and_hms_opt(0, 0, 0)? + chrono::Duration::weeks(1);
Some(next_naive.and_utc())
}
"monthly" => {
let next_naive =
(now.date_naive() + chrono::Duration::days(1)).and_hms_opt(0, 0, 0)?;
Some(next_naive.and_utc())
}
_ => {
let next_naive = now.date_naive().and_hms_opt(0, 0, 0)? + chrono::Duration::days(1);
Some(next_naive.and_utc())
}
}
}
fn should_rotate_by_time_inner(&self, inner: &FileSinkInner) -> bool {
let now = Utc::now();
let current_date = now.date_naive().num_days_from_ce();
if self.config.rotation_time == "daily" || self.config.rotation_time == "weekly" {
if let Some(last_date) = inner.last_rotation_date {
if current_date > last_date {
return true;
}
}
}
if let Some(next_time) = inner.next_rotation_time {
if now >= next_time {
return true;
}
}
false
}
fn update_next_rotation_time_inner(&self, inner: &mut FileSinkInner) {
inner.next_rotation_time = Self::calculate_next_rotation_time(&self.config.rotation_time);
}
fn start_rotation_timer(&self) {
let rotation_interval = self.rotation_interval;
let last_rotation;
{
let inner = self.inner.read();
last_rotation = Arc::new(parking_lot::Mutex::new(inner.last_rotation));
}
{
let mut inner = self.inner.write();
inner.rotation_timer = Some(last_rotation.clone());
}
let shutdown_flag = self.shutdown_flag.clone();
let timer_handle = thread::spawn(move || {
let check_interval = StdDuration::from_secs(60); loop {
if shutdown_flag.load(Ordering::Relaxed) {
break;
}
let mut elapsed = StdDuration::ZERO;
const POLL_INTERVAL: StdDuration = StdDuration::from_millis(100);
while elapsed < check_interval {
if shutdown_flag.load(Ordering::Relaxed) {
break;
}
let step = std::cmp::min(POLL_INTERVAL, check_interval - elapsed);
thread::sleep(step);
elapsed += step;
}
if shutdown_flag.load(Ordering::Relaxed) {
break;
}
let mut last_rotation_guard = last_rotation.lock();
if last_rotation_guard.elapsed() >= rotation_interval {
*last_rotation_guard =
Instant::now() - rotation_interval + StdDuration::from_secs(1);
}
}
});
self.inner.write().timer_handle = Some(timer_handle);
}
fn flush_batch_inner(&self, inner: &mut FileSinkInner) -> Result<(), InklogError> {
if inner.batch_buffer.is_empty() {
return Ok(());
}
let records = std::mem::take(&mut inner.batch_buffer);
if let Some(file) = &mut inner.current_file {
for record in &records {
match writeln!(
file,
"{} [{}] {} - {}",
record.timestamp.to_rfc3339(),
record.level,
record.target,
record.message
) {
Ok(_) => {
let len = record.timestamp.to_rfc3339().len()
+ record.level.len()
+ record.target.len()
+ record.message.len()
+ 7;
inner.current_size += len as u64;
}
Err(e) => {
error!("Batch write error: {}", e);
inner.circuit_breaker.record_failure();
let _ = self.open_file_inner(inner);
break;
}
}
}
inner.circuit_breaker.record_success();
}
inner.last_flush_time = Instant::now();
self.check_rotation_inner(inner)?;
Ok(())
}
fn compress_file(&self, path: &PathBuf) -> Result<PathBuf, InklogError> {
let compressed_path = path.with_extension("zst");
let input_file = fs::File::open(path).map_err(|e| {
error!("Failed to open file for compression: {}", e);
InklogError::IoError(e)
})?;
let output_file = fs::File::create(&compressed_path).map_err(|e| {
error!("Failed to create compressed file: {}", e);
InklogError::IoError(e)
})?;
let mut encoder = zstd::stream::Encoder::new(output_file, self.config.compression_level)
.map_err(|e| InklogError::CompressionError(e.to_string()))?
.auto_finish();
let mut reader = std::io::BufReader::new(input_file);
let mut buffer = [0u8; 8192];
loop {
let bytes_read = std::io::Read::read(&mut reader, &mut buffer)?;
if bytes_read == 0 {
break;
}
std::io::Write::write_all(&mut encoder, &buffer[..bytes_read])?;
}
drop(encoder);
if self.config.encrypt {
let encrypted_path = compressed_path.with_extension("zst.enc");
if let Err(e) = self.encrypt_file(&compressed_path, &encrypted_path) {
error!("Encryption failed: {}", e);
let _ = fs::rename(
&compressed_path,
encrypted_path.with_extension("zst.unencrypted"),
);
return Err(e);
}
let _ = fs::remove_file(&compressed_path);
Ok(encrypted_path)
} else {
let _ = fs::remove_file(path);
Ok(compressed_path)
}
}
fn encrypt_file(&self, input_path: &PathBuf, output_path: &PathBuf) -> Result<(), InklogError> {
use aes_gcm::{Aes256Gcm, Nonce};
use rand::Rng;
let key_bytes = self.get_encryption_key()?;
let cipher = Aes256Gcm::new_from_slice(&key_bytes)
.map_err(|e| InklogError::EncryptionError(format!("Invalid key: {}", e)))?;
let mut nonce_bytes = [0u8; 12];
rand::rng().fill_bytes(&mut nonce_bytes);
let nonce = Nonce::from(nonce_bytes);
let input_data = fs::read(input_path).map_err(|e| {
error!("Failed to read file for encryption: {}", e);
InklogError::IoError(e)
})?;
let ciphertext = cipher.encrypt(&nonce, input_data.as_slice()).map_err(|e| {
error!("Encryption failed: {}", e);
InklogError::EncryptionError(e.to_string())
})?;
let mut output = fs::File::create(output_path).map_err(|e| {
error!("Failed to create encrypted file: {}", e);
InklogError::IoError(e)
})?;
output.write_all(&nonce_bytes)?;
output.write_all(&ciphertext)?;
debug!("Encrypted log file: {}", output_path.display());
Ok(())
}
fn rotate_inner(&self, inner: &mut FileSinkInner) -> Result<(), InklogError> {
debug!("Rotating log file: {}", self.config.path.display());
let _ = inner.current_file.take();
let timestamp = chrono::Utc::now().format("%Y%m%d_%H%M%S");
let new_path = if let Some(parent) = self.config.path.parent() {
let stem = self.config.path.file_stem().unwrap_or_default();
let ext = self.config.path.extension().unwrap_or_default();
parent.join(format!(
"{}_{}.{}",
stem.to_string_lossy(),
timestamp,
ext.to_string_lossy()
))
} else {
PathBuf::from(format!("{}_{}", self.config.path.display(), timestamp))
};
if self.config.path.exists() {
if let Err(e) = fs::rename(&self.config.path, &new_path) {
error!("Failed to rename log file: {}", e);
if fs::copy(&self.config.path, &new_path).is_ok() {
let _ = fs::remove_file(&self.config.path);
} else {
return Err(InklogError::IoError(e));
}
}
}
inner.sequence += 1;
inner.last_rotation = Instant::now();
self.update_next_rotation_time_inner(inner);
inner.current_size = 0;
info!("Log rotated to: {}", new_path.display());
if self.config.compress {
let config = self.config.clone();
let path = new_path.clone();
let _ = thread::spawn(move || {
let inner = FileSinkInner {
current_file: None,
current_size: 0,
last_rotation: Instant::now(),
next_rotation_time: None,
last_rotation_date: None,
sequence: 0,
fallback_sink: None,
circuit_breaker: CircuitBreaker::new(5, StdDuration::from_secs(30), 3),
batch_buffer: Vec::new(),
last_flush_time: Instant::now(),
timer_handle: None,
rotation_timer: None,
cleanup_timer_handle: None,
rotation_strategy: Box::new(
crate::support::io::sink::rotation::CompositeRotation::new(vec![]),
),
};
let sink = FileSink {
config,
rotation_interval: StdDuration::from_secs(86400),
last_cleanup_time: Arc::new(parking_lot::Mutex::new(None)),
shutdown_flag: Arc::new(AtomicBool::new(false)),
masker: DataMasker::new(),
inner: RwLock::new(inner),
};
if let Err(e) = sink.compress_file(&path) {
error!("Failed to compress rotated log: {}", e);
}
});
} else if self.config.encrypt {
let config = self.config.clone();
let path = new_path.clone();
let _ = thread::spawn(move || {
let inner = FileSinkInner {
current_file: None,
current_size: 0,
last_rotation: Instant::now(),
next_rotation_time: None,
last_rotation_date: None,
sequence: 0,
fallback_sink: None,
circuit_breaker: CircuitBreaker::new(5, StdDuration::from_secs(30), 3),
batch_buffer: Vec::new(),
last_flush_time: Instant::now(),
timer_handle: None,
rotation_timer: None,
cleanup_timer_handle: None,
rotation_strategy: Box::new(
crate::support::io::sink::rotation::CompositeRotation::new(vec![]),
),
};
let sink = FileSink {
config,
rotation_interval: StdDuration::from_secs(86400),
last_cleanup_time: Arc::new(parking_lot::Mutex::new(None)),
shutdown_flag: Arc::new(AtomicBool::new(false)),
masker: DataMasker::new(),
inner: RwLock::new(inner),
};
let encrypted_path = path.with_extension("enc");
if let Err(e) = sink.encrypt_file(&path, &encrypted_path) {
error!("Failed to encrypt rotated log: {}", e);
} else {
let _ = fs::remove_file(&path);
}
});
}
self.open_file_inner(inner)
}
fn check_rotation_inner(&self, inner: &mut FileSinkInner) -> Result<(), InklogError> {
let rotate_by_size =
Self::parse_size(&self.config.max_size).is_some_and(|max| inner.current_size >= max);
let rotate_by_time = self.should_rotate_by_time_inner(inner);
if rotate_by_size || rotate_by_time {
self.rotate_inner(inner)?;
}
Ok(())
}
}
#[async_trait]
impl LogSink for FileSink {
async fn write(&self, record: &LogRecord) -> Result<(), InklogError> {
let circuit_open = {
let inner = self.inner.read();
!inner.circuit_breaker.can_execute()
};
if circuit_open {
let fallback = self.inner.read().fallback_sink.clone();
if let Some(sink) = fallback {
let _ = sink.write(record).await;
}
return Ok(());
}
if !self.check_disk_space()? {
warn!("Low disk space - checking before write");
let fallback = self.inner.read().fallback_sink.clone();
if let Some(sink) = fallback {
let _ = sink.write(record).await;
}
return Ok(());
}
let rotation_failed_fallback: Option<Arc<dyn LogSink + Send + Sync>> = {
let mut inner = self.inner.write();
let masked_record = if self.config.masking_enabled {
let mut masked = record.clone();
masked.message = self.masker.mask(&record.message);
self.masker.mask_hashmap(&mut masked.fields);
masked
} else {
record.clone()
};
let record_len = masked_record.timestamp.to_rfc3339().len()
+ masked_record.level.len()
+ masked_record.target.len()
+ masked_record.message.len()
+ 7;
inner.current_size += record_len as u64;
inner.batch_buffer.push(masked_record);
let should_rotate = Self::parse_size(&self.config.max_size)
.is_some_and(|max| inner.current_size >= max)
|| inner
.rotation_timer
.as_ref()
.map(|t| t.lock().elapsed() >= self.rotation_interval)
.unwrap_or(false);
if should_rotate {
if let Err(e) = self.rotate_inner(&mut inner) {
error!("Rotation failed: {}", e);
inner.fallback_sink.clone()
} else {
let now = Instant::now();
let flush_interval = StdDuration::from_millis(self.config.flush_interval_ms);
if inner.batch_buffer.len() >= self.config.batch_size
|| now.duration_since(inner.last_flush_time) >= flush_interval
{
self.flush_batch_inner(&mut inner)?;
}
None
}
} else {
let now = Instant::now();
let flush_interval = StdDuration::from_millis(self.config.flush_interval_ms);
if inner.batch_buffer.len() >= self.config.batch_size
|| now.duration_since(inner.last_flush_time) >= flush_interval
{
self.flush_batch_inner(&mut inner)?;
}
None
}
};
if let Some(sink) = rotation_failed_fallback {
let _ = sink.write(record).await;
}
Ok(())
}
async fn flush(&self) -> Result<(), InklogError> {
let mut inner = self.inner.write();
self.flush_batch_inner(&mut inner)?;
if let Some(file) = &mut inner.current_file {
file.flush()?;
}
Ok(())
}
fn is_healthy(&self) -> bool {
self.inner.read().current_file.is_some()
}
async fn shutdown(&self) -> Result<(), InklogError> {
self.shutdown_flag.store(true, Ordering::Relaxed);
let fallback = {
let mut inner = self.inner.write();
if let Some(handle) = inner.timer_handle.take() {
let _ = handle.join();
}
inner.rotation_timer = None;
if let Some(handle) = inner.cleanup_timer_handle.take() {
let _ = handle.join();
}
self.flush_batch_inner(&mut inner)?;
if let Some(file) = &mut inner.current_file {
file.flush()?;
}
inner.fallback_sink.take()
};
if let Some(sink) = fallback {
let _ = sink.shutdown().await;
}
Ok(())
}
}
impl Drop for FileSink {
fn drop(&mut self) {
const SHUTDOWN_TIMEOUT_MS: u64 = 5000;
self.shutdown_flag.store(true, Ordering::SeqCst);
{
let mut inner = self.inner.write();
let _ = self.flush_batch_inner(&mut inner);
if let Some(mut file) = inner.current_file.take() {
let _ = file.flush();
}
}
{
let mut inner = self.inner.write();
if let Some(handle) = inner.timer_handle.take() {
let start = std::time::Instant::now();
while !handle.is_finished() {
if start.elapsed().as_millis() > SHUTDOWN_TIMEOUT_MS as u128 {
tracing::warn!(
"Warning: rotation timer shutdown timeout after {}ms",
SHUTDOWN_TIMEOUT_MS
);
break;
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
}
}
{
let mut inner = self.inner.write();
if let Some(handle) = inner.cleanup_timer_handle.take() {
let start = std::time::Instant::now();
while !handle.is_finished() {
if start.elapsed().as_millis() > SHUTDOWN_TIMEOUT_MS as u128 {
tracing::warn!(
"Warning: cleanup timer shutdown timeout after {}ms",
SHUTDOWN_TIMEOUT_MS
);
break;
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
}
}
{
let mut inner = self.inner.write();
let _fallback = inner.fallback_sink.take();
}
}
}
impl std::fmt::Debug for FileSink {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let inner = self.inner.read();
f.debug_struct("FileSink")
.field("path", &self.config.path)
.field("current_size", &inner.current_size)
.field("circuit_breaker", &inner.circuit_breaker)
.finish()
}
}
impl Clone for FileSink {
fn clone(&self) -> Self {
let inner = FileSinkInner {
current_file: None,
current_size: 0,
last_rotation: Instant::now(),
next_rotation_time: None,
last_rotation_date: None,
sequence: 0,
fallback_sink: None,
circuit_breaker: CircuitBreaker::new(5, StdDuration::from_secs(30), 3),
batch_buffer: Vec::with_capacity(self.config.batch_size),
last_flush_time: Instant::now(),
timer_handle: None,
rotation_timer: None,
cleanup_timer_handle: None,
rotation_strategy: self.inner.read().rotation_strategy.clone_boxed(),
};
Self {
config: self.config.clone(),
rotation_interval: self.rotation_interval,
last_cleanup_time: Arc::new(parking_lot::Mutex::new(None)),
shutdown_flag: Arc::new(AtomicBool::new(false)),
masker: DataMasker::new(),
inner: RwLock::new(inner),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::FileSinkConfig;
use crate::LogRecord;
use base64::Engine;
use chrono::Timelike;
use chrono::Utc;
use serial_test::serial;
use std::collections::HashMap;
use tempfile::tempdir;
#[allow(dead_code)]
fn create_test_record(message: &str) -> LogRecord {
LogRecord {
timestamp: Utc::now(),
level: "INFO".to_string(),
target: "test_module".to_string(),
message: message.to_string(),
fields: HashMap::new(),
file: Some("/path/to/test.rs".to_string()),
line: Some(42),
thread_id: "test-thread".to_string(),
}
}
fn create_test_file_sink(config: FileSinkConfig) -> FileSink {
let inner = FileSinkInner {
current_file: None,
current_size: 0,
last_rotation: Instant::now(),
next_rotation_time: None,
last_rotation_date: None,
sequence: 0,
fallback_sink: None,
circuit_breaker: CircuitBreaker::new(5, StdDuration::from_secs(30), 3),
batch_buffer: Vec::new(),
last_flush_time: Instant::now(),
timer_handle: None,
rotation_timer: None,
cleanup_timer_handle: None,
rotation_strategy: Box::new(
crate::support::io::sink::rotation::CompositeRotation::new(vec![]),
),
};
FileSink {
config,
rotation_interval: StdDuration::from_secs(86400),
last_cleanup_time: Arc::new(parking_lot::Mutex::new(None)),
shutdown_flag: Arc::new(AtomicBool::new(false)),
masker: DataMasker::new(),
inner: RwLock::new(inner),
}
}
#[test]
fn test_parse_size() {
assert_eq!(FileSink::parse_size("100"), Some(100));
assert_eq!(FileSink::parse_size("100KB"), Some(100 * 1024));
assert_eq!(FileSink::parse_size("10MB"), Some(10 * 1024 * 1024));
assert_eq!(FileSink::parse_size("1GB"), Some(1024 * 1024 * 1024));
assert_eq!(FileSink::parse_size(" 5MB "), Some(5 * 1024 * 1024));
assert_eq!(FileSink::parse_size("invalid"), None);
}
#[test]
fn test_perform_cleanup() {
let dir = tempdir().unwrap();
let log_path = dir.path().join("test.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
max_size: "1MB".to_string(),
rotation_time: "daily".to_string(),
keep_files: 2,
compress: false,
compression_level: 3,
encrypt: false,
encryption_key_env: None,
retention_days: 30,
max_total_size: "1GB".to_string(),
cleanup_interval_minutes: 60,
batch_size: 100,
flush_interval_ms: 100,
masking_enabled: true,
};
let old_file = dir.path().join("test_old.log");
std::fs::write(&old_file, "old content").unwrap();
let result = FileSink::perform_cleanup(&config, &log_path);
assert!(result.is_ok());
}
#[test]
fn test_get_encryption_key() {
let config = FileSinkConfig {
enabled: true,
path: PathBuf::from("test.log"),
encryption_key_env: Some("TEST_KEY".to_string()),
..Default::default()
};
std::env::set_var("TEST_KEY", "YWJjZGVmZ2hpamtsbW5vcHFyc3R1dnd4eXoxMjM0NTY=");
let sink = create_test_file_sink(config);
let key_result = sink.get_encryption_key();
assert!(key_result.is_ok());
assert_eq!(key_result.unwrap().len(), 32);
std::env::remove_var("TEST_KEY");
}
#[test]
fn test_disk_space_info() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.get_disk_space_info();
assert!(result.is_ok());
let (total, available) = result.unwrap();
assert!(total > 0);
assert!(available > 0);
}
#[test]
fn test_check_disk_space_logic() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.check_disk_space();
assert!(result.is_ok());
}
#[tokio::test]
async fn test_write_with_disk_space_check() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("test.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
..Default::default()
};
let sink = FileSink::new(config).unwrap();
let record = LogRecord {
timestamp: chrono::Utc::now(),
level: "INFO".to_string(),
target: "test".to_string(),
message: "Test message".to_string(),
fields: HashMap::new(),
file: Some("test.rs".to_string()),
line: Some(1),
thread_id: format!("{:?}", std::thread::current().id()),
};
let result = sink.write(&record).await;
assert!(
result.is_ok(),
"Write should succeed with sufficient disk space"
);
sink.flush().await.unwrap();
assert!(log_path.exists(), "Log file should exist");
}
#[test]
fn test_parse_size_kb() {
assert_eq!(FileSink::parse_size("500KB"), Some(500 * 1024));
}
#[test]
fn test_parse_size_mb() {
assert_eq!(FileSink::parse_size("2MB"), Some(2 * 1024 * 1024));
}
#[test]
fn test_parse_size_gb() {
assert_eq!(FileSink::parse_size("1GB"), Some(1024 * 1024 * 1024));
}
#[test]
fn test_parse_size_with_spaces() {
assert_eq!(FileSink::parse_size(" 3MB "), Some(3 * 1024 * 1024));
}
#[test]
fn test_parse_size_invalid() {
assert_eq!(FileSink::parse_size("invalid"), None);
assert_eq!(FileSink::parse_size(""), None);
}
#[test]
fn test_parse_size_zero() {
assert_eq!(FileSink::parse_size("0"), Some(0));
assert_eq!(FileSink::parse_size("0MB"), Some(0));
}
#[test]
fn test_get_encryption_key_missing_env() {
let config = FileSinkConfig {
enabled: true,
path: PathBuf::from("test.log"),
encryption_key_env: Some("MISSING_KEY".to_string()),
..Default::default()
};
std::env::remove_var("MISSING_KEY");
let sink = create_test_file_sink(config);
let result = sink.get_encryption_key();
assert!(result.is_err());
}
#[test]
fn test_get_encryption_key_no_env_var() {
let config = FileSinkConfig {
enabled: true,
path: PathBuf::from("test.log"),
encryption_key_env: None,
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.get_encryption_key();
if result.is_err() {
assert!(std::env::var("LOG_ENCRYPTION_KEY").is_err());
}
}
#[test]
fn test_file_sink_new_default() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
println!("FileSinkConfig: {:?}", config);
let result = FileSink::new(config);
if let Err(ref e) = result {
println!("Error: {:?}", e);
}
assert!(
result.is_ok(),
"Expected FileSink::new to succeed with default config, but got error: {:?}",
result.err()
);
}
#[test]
fn test_file_sink_new_with_path() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
let result = FileSink::new(config);
assert!(result.is_ok());
}
#[test]
fn test_file_sink_disabled() {
let config = FileSinkConfig {
enabled: false,
path: PathBuf::from("test.log"),
..Default::default()
};
let result = FileSink::new(config);
assert!(result.is_ok());
}
#[test]
fn test_validate_key_entropy_strong() {
let strong_key = [
0x3a, 0x7b, 0x9c, 0x1d, 0x4e, 0x8f, 0x2c, 0x6b, 0x9a, 0x3d, 0x8e, 0x1f, 0x4a, 0x7d,
0x2e, 0x6f, 0x9b, 0x3c, 0x8d, 0x1e, 0x4b, 0x6a, 0x2b, 0x6c, 0x9f, 0x3a, 0x8b, 0x1c,
0x4d, 0x7e, 0x2f, 0x6a,
];
assert!(FileSink::validate_key_entropy(&strong_key).is_ok());
}
#[test]
fn test_validate_key_entropy_weak() {
let weak_key = [0xaa; 32];
assert!(FileSink::validate_key_entropy(&weak_key).is_err());
}
#[test]
fn test_validate_key_entropy_empty() {
let empty_key: [u8; 0] = [];
assert!(FileSink::validate_key_entropy(&empty_key).is_err());
}
#[test]
fn test_get_encryption_key_too_short() {
let config = FileSinkConfig {
enabled: true,
path: PathBuf::from("test.log"),
encryption_key_env: Some("TEST_SHORT_KEY".to_string()),
..Default::default()
};
std::env::set_var("TEST_SHORT_KEY", "YWJjZA==");
let sink = create_test_file_sink(config);
let result = sink.get_encryption_key();
assert!(result.is_err());
assert!(result
.unwrap_err()
.to_string()
.contains("at least 16 characters"));
}
#[test]
fn test_nonce_generation_unique() {
use rand::Rng;
let mut nonces = Vec::new();
for _ in 0..100 {
let mut nonce_bytes = [0u8; 12];
rand::rng().fill_bytes(&mut nonce_bytes);
nonces.push(nonce_bytes);
}
for i in 0..nonces.len() {
for j in (i + 1)..nonces.len() {
assert_ne!(
nonces[i], nonces[j],
"Nonce {} and {} should be different",
i, j
);
}
}
}
#[allow(dead_code)]
fn make_test_key() -> (Vec<u8>, String) {
let key_bytes: Vec<u8> = vec![
0x3a, 0x7b, 0x9c, 0x1d, 0x4e, 0x8f, 0x2c, 0x6b, 0x9a, 0x3d, 0x8e, 0x1f, 0x4a, 0x7d,
0x2e, 0x6f, 0x9b, 0x3c, 0x8d, 0x1e, 0x4b, 0x6a, 0x2b, 0x6c, 0x9f, 0x3a, 0x8b, 0x1c,
0x4d, 0x7e, 0x2f, 0x6a,
];
let key_b64 = base64::engine::general_purpose::STANDARD.encode(&key_bytes);
(key_bytes, key_b64)
}
#[test]
fn test_parse_size_tb() {
assert_eq!(FileSink::parse_size("1TB"), Some(1024 * 1024 * 1024 * 1024));
assert_eq!(
FileSink::parse_size("2TB"),
Some(2 * 1024 * 1024 * 1024 * 1024)
);
}
#[test]
fn test_parse_size_decimal_rejected() {
assert_eq!(FileSink::parse_size("1.5MB"), None);
assert_eq!(FileSink::parse_size("0.5"), None);
}
#[test]
fn test_parse_size_negative_rejected() {
assert_eq!(FileSink::parse_size("-100"), None);
}
#[test]
fn test_calculate_next_rotation_time_hourly() {
let now = Utc::now();
let result = FileSink::calculate_next_rotation_time("hourly");
assert!(result.is_some());
let next = result.unwrap();
assert!(next > now);
let diff = next - now;
assert!(
diff.num_minutes() >= 59 && diff.num_minutes() <= 61,
"hourly rotation should be ~60 minutes away, got {}",
diff.num_minutes()
);
}
#[test]
fn test_calculate_next_rotation_time_daily() {
let now = Utc::now();
let result = FileSink::calculate_next_rotation_time("daily");
assert!(result.is_some());
let next = result.unwrap();
assert_eq!(next.hour(), 0);
assert_eq!(next.minute(), 0);
assert_eq!(next.second(), 0);
assert!(next > now);
}
#[test]
fn test_calculate_next_rotation_time_weekly() {
let now = Utc::now();
let result = FileSink::calculate_next_rotation_time("weekly");
assert!(result.is_some());
let next = result.unwrap();
assert_eq!(next.hour(), 0);
assert_eq!(next.minute(), 0);
assert_eq!(next.second(), 0);
assert!(next > now);
}
#[test]
fn test_calculate_next_rotation_time_monthly() {
let result = FileSink::calculate_next_rotation_time("monthly");
assert!(result.is_some());
}
#[test]
fn test_calculate_next_rotation_time_invalid_defaults_to_daily() {
let result = FileSink::calculate_next_rotation_time("invalid_interval");
assert!(result.is_some());
let next = result.unwrap();
assert_eq!(next.hour(), 0);
assert_eq!(next.minute(), 0);
assert_eq!(next.second(), 0);
}
#[test]
fn test_update_next_rotation_time_inner_sets_value() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
rotation_time: "hourly".to_string(),
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
inner.next_rotation_time = None;
sink.update_next_rotation_time_inner(&mut inner);
assert!(inner.next_rotation_time.is_some());
}
#[test]
fn test_should_rotate_by_time_inner_no_next_time() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
rotation_time: "daily".to_string(),
..Default::default()
};
let sink = create_test_file_sink(config);
let inner = sink.inner.read();
let result = sink.should_rotate_by_time_inner(&inner);
assert!(!result);
}
#[test]
fn test_should_rotate_by_time_inner_past_next_time() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
rotation_time: "daily".to_string(),
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
inner.next_rotation_time = Some(Utc::now() - chrono::Duration::hours(1));
let result = sink.should_rotate_by_time_inner(&inner);
assert!(result);
}
#[test]
fn test_should_rotate_by_time_inner_future_next_time() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
rotation_time: "daily".to_string(),
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
inner.next_rotation_time = Some(Utc::now() + chrono::Duration::hours(1));
let result = sink.should_rotate_by_time_inner(&inner);
assert!(!result);
}
#[test]
fn test_should_rotate_by_time_inner_daily_date_change() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
rotation_time: "daily".to_string(),
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
let yesterday = Utc::now().date_naive().num_days_from_ce() - 1;
inner.last_rotation_date = Some(yesterday);
inner.next_rotation_time = Some(Utc::now() + chrono::Duration::days(1));
let result = sink.should_rotate_by_time_inner(&inner);
assert!(result);
}
#[test]
fn test_open_file_inner_creates_nested_directory() {
let temp_dir = tempdir().unwrap();
let nested = temp_dir.path().join("nested").join("deep");
let log_path = nested.join("test.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
let result = sink.open_file_inner(&mut inner);
assert!(result.is_ok());
assert!(inner.current_file.is_some());
assert!(log_path.exists());
}
#[test]
fn test_open_file_inner_detects_existing_size() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("test.log");
let existing = "existing content\n";
std::fs::write(&log_path, existing).unwrap();
let config = FileSinkConfig {
enabled: true,
path: log_path,
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
let result = sink.open_file_inner(&mut inner);
assert!(result.is_ok());
assert_eq!(inner.current_size, existing.len() as u64);
}
#[test]
fn test_flush_batch_inner_empty_buffer_noop() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
sink.open_file_inner(&mut inner).unwrap();
let result = sink.flush_batch_inner(&mut inner);
assert!(result.is_ok());
}
#[test]
fn test_flush_batch_inner_writes_records_to_file() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("test.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
sink.open_file_inner(&mut inner).unwrap();
inner.batch_buffer.push(create_test_record("Message 1"));
inner.batch_buffer.push(create_test_record("Message 2"));
let result = sink.flush_batch_inner(&mut inner);
assert!(result.is_ok());
assert!(inner.batch_buffer.is_empty());
drop(inner);
let content = std::fs::read_to_string(&log_path).unwrap();
assert!(content.contains("Message 1"));
assert!(content.contains("Message 2"));
}
#[test]
fn test_flush_batch_inner_increments_current_size() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
sink.open_file_inner(&mut inner).unwrap();
let initial_size = inner.current_size;
inner.batch_buffer.push(create_test_record("Test message"));
sink.flush_batch_inner(&mut inner).unwrap();
assert!(inner.current_size > initial_size);
}
#[test]
fn test_check_rotation_inner_no_rotation_needed() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
max_size: "1MB".to_string(),
rotation_time: "daily".to_string(),
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
sink.open_file_inner(&mut inner).unwrap();
inner.current_size = 100; inner.next_rotation_time = Some(Utc::now() + chrono::Duration::days(1));
let result = sink.check_rotation_inner(&mut inner);
assert!(result.is_ok());
assert_eq!(inner.sequence, 0); }
#[test]
fn test_check_rotation_inner_by_size_triggers_rotation() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
max_size: "100".to_string(), rotation_time: "daily".to_string(),
compress: false,
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
sink.open_file_inner(&mut inner).unwrap();
std::fs::write(sink.config.path.clone(), "x").unwrap();
inner.current_size = 200; inner.next_rotation_time = Some(Utc::now() + chrono::Duration::days(1));
let result = sink.check_rotation_inner(&mut inner);
assert!(result.is_ok());
assert_eq!(inner.sequence, 1); }
#[test]
fn test_rotate_inner_renames_original_file() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("test.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
compress: false,
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
sink.open_file_inner(&mut inner).unwrap();
std::fs::write(&log_path, "test content").unwrap();
let result = sink.rotate_inner(&mut inner);
assert!(result.is_ok());
assert!(log_path.exists());
let count = std::fs::read_dir(temp_dir.path()).unwrap().count();
assert!(count >= 2);
}
#[test]
fn test_rotate_inner_increments_sequence() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
compress: false,
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
sink.open_file_inner(&mut inner).unwrap();
let initial = inner.sequence;
sink.rotate_inner(&mut inner).unwrap();
assert_eq!(inner.sequence, initial + 1);
sink.rotate_inner(&mut inner).unwrap();
assert_eq!(inner.sequence, initial + 2);
}
#[test]
fn test_rotate_inner_resets_current_size() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
compress: false,
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
sink.open_file_inner(&mut inner).unwrap();
inner.current_size = 5000;
sink.rotate_inner(&mut inner).unwrap();
assert_eq!(inner.current_size, 0);
}
#[test]
fn test_rotate_inner_updates_next_rotation_time() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
rotation_time: "hourly".to_string(),
compress: false,
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
sink.open_file_inner(&mut inner).unwrap();
inner.next_rotation_time = None;
sink.rotate_inner(&mut inner).unwrap();
assert!(inner.next_rotation_time.is_some());
}
#[test]
fn test_compress_file_roundtrip() {
let temp_dir = tempdir().unwrap();
let original_path = temp_dir.path().join("test.log");
let original_content = b"This is test content for compression. Hello World!";
std::fs::write(&original_path, original_content).unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("dummy.log"),
compress: true,
compression_level: 3,
encrypt: false,
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.compress_file(&original_path);
assert!(result.is_ok());
let compressed_path = result.unwrap();
assert_eq!(compressed_path.extension().unwrap(), "zst");
assert!(compressed_path.exists());
assert!(!original_path.exists());
let compressed_file = std::fs::File::open(&compressed_path).unwrap();
let mut decoder = zstd::stream::Decoder::new(compressed_file).unwrap();
let mut decompressed = Vec::new();
std::io::Read::read_to_end(&mut decoder, &mut decompressed).unwrap();
assert_eq!(decompressed, original_content);
}
#[test]
fn test_compress_file_nonexistent_input_returns_error() {
let temp_dir = tempdir().unwrap();
let nonexistent = temp_dir.path().join("nonexistent.log");
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("dummy.log"),
compress: true,
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.compress_file(&nonexistent);
assert!(result.is_err());
}
#[test]
#[serial]
fn test_compress_file_with_encryption_roundtrip() {
let temp_dir = tempdir().unwrap();
let original_path = temp_dir.path().join("test.log");
let original_content = b"Sensitive log content that needs encryption";
std::fs::write(&original_path, original_content).unwrap();
let (key_bytes, key_b64) = make_test_key();
std::env::set_var("TEST_COMPRESS_ENC_KEY", &key_b64);
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("dummy.log"),
compress: true,
compression_level: 3,
encrypt: true,
encryption_key_env: Some("TEST_COMPRESS_ENC_KEY".to_string()),
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.compress_file(&original_path);
assert!(result.is_ok(), "compress_file failed: {:?}", result.err());
let encrypted_path = result.unwrap();
assert_eq!(encrypted_path.extension().unwrap(), "enc");
assert!(encrypted_path.exists());
let encrypted_data = std::fs::read(&encrypted_path).unwrap();
assert!(encrypted_data.len() > 12);
use aes_gcm::{Aes256Gcm, Nonce};
let cipher = Aes256Gcm::new_from_slice(&key_bytes).unwrap();
let nonce_arr: [u8; 12] = encrypted_data[..12].try_into().unwrap();
let nonce = Nonce::from(nonce_arr);
let ciphertext = &encrypted_data[12..];
let decrypted_compressed = cipher.decrypt(&nonce, ciphertext).unwrap();
let mut decoder = zstd::stream::Decoder::new(&decrypted_compressed[..]).unwrap();
let mut decompressed = Vec::new();
std::io::Read::read_to_end(&mut decoder, &mut decompressed).unwrap();
assert_eq!(decompressed, original_content);
std::env::remove_var("TEST_COMPRESS_ENC_KEY");
}
#[test]
#[serial]
fn test_encrypt_file_roundtrip() {
let temp_dir = tempdir().unwrap();
let input_path = temp_dir.path().join("test.log");
let output_path = temp_dir.path().join("test.log.enc");
let original_content = b"Secret log content for encryption test";
std::fs::write(&input_path, original_content).unwrap();
let (key_bytes, key_b64) = make_test_key();
std::env::set_var("TEST_ENC_KEY_RT", &key_b64);
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("dummy.log"),
encrypt: true,
encryption_key_env: Some("TEST_ENC_KEY_RT".to_string()),
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.encrypt_file(&input_path, &output_path);
assert!(result.is_ok(), "encrypt_file failed: {:?}", result.err());
assert!(output_path.exists());
let encrypted_data = std::fs::read(&output_path).unwrap();
assert!(encrypted_data.len() > 12);
use aes_gcm::{Aes256Gcm, Nonce};
let cipher = Aes256Gcm::new_from_slice(&key_bytes).unwrap();
let nonce_arr: [u8; 12] = encrypted_data[..12].try_into().unwrap();
let nonce = Nonce::from(nonce_arr);
let ciphertext = &encrypted_data[12..];
let decrypted = cipher.decrypt(&nonce, ciphertext).unwrap();
assert_eq!(decrypted, original_content);
std::env::remove_var("TEST_ENC_KEY_RT");
}
#[test]
#[serial]
fn test_encrypt_file_missing_key_returns_error() {
let temp_dir = tempdir().unwrap();
let input_path = temp_dir.path().join("input.log");
let output_path = temp_dir.path().join("output.log.enc");
std::fs::write(&input_path, "content").unwrap();
std::env::remove_var("TEST_MISSING_ENC_KEY_VAR");
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("dummy.log"),
encrypt: true,
encryption_key_env: Some("TEST_MISSING_ENC_KEY_VAR".to_string()),
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.encrypt_file(&input_path, &output_path);
assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("not found"));
}
#[test]
#[serial]
fn test_encrypt_file_nonexistent_input_returns_error() {
let temp_dir = tempdir().unwrap();
let input_path = temp_dir.path().join("nonexistent.log");
let output_path = temp_dir.path().join("output.log.enc");
let (_key_bytes, key_b64) = make_test_key();
std::env::set_var("TEST_ENC_KEY_NI", &key_b64);
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("dummy.log"),
encrypt: true,
encryption_key_env: Some("TEST_ENC_KEY_NI".to_string()),
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.encrypt_file(&input_path, &output_path);
assert!(result.is_err());
std::env::remove_var("TEST_ENC_KEY_NI");
}
#[test]
#[serial]
fn test_encrypt_file_invalid_base64_key_returns_error() {
let temp_dir = tempdir().unwrap();
let input_path = temp_dir.path().join("input.log");
let output_path = temp_dir.path().join("output.log.enc");
std::fs::write(&input_path, "content").unwrap();
std::env::set_var("TEST_INVALID_B64_KEY", "not_valid_base64!!!*@$");
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("dummy.log"),
encrypt: true,
encryption_key_env: Some("TEST_INVALID_B64_KEY".to_string()),
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.encrypt_file(&input_path, &output_path);
assert!(result.is_err());
std::env::remove_var("TEST_INVALID_B64_KEY");
}
#[test]
#[serial]
fn test_encrypt_file_wrong_length_key_returns_error() {
let temp_dir = tempdir().unwrap();
let input_path = temp_dir.path().join("input.log");
let output_path = temp_dir.path().join("output.log.enc");
std::fs::write(&input_path, "content").unwrap();
let short_key = base64::engine::general_purpose::STANDARD.encode(b"1234567890123456");
std::env::set_var("TEST_WRONG_LEN_KEY", &short_key);
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("dummy.log"),
encrypt: true,
encryption_key_env: Some("TEST_WRONG_LEN_KEY".to_string()),
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.encrypt_file(&input_path, &output_path);
assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("32 bytes"));
std::env::remove_var("TEST_WRONG_LEN_KEY");
}
#[test]
fn test_perform_cleanup_removes_excess_by_total_size() {
let dir = tempdir().unwrap();
let log_path = dir.path().join("test.log");
for i in 0..5 {
let p = dir.path().join(format!("test_{}.log", i));
std::fs::write(&p, "x".repeat(1024)).unwrap();
}
let config = FileSinkConfig {
enabled: true,
path: log_path,
max_size: "1MB".to_string(),
rotation_time: "daily".to_string(),
keep_files: 2,
compress: false,
compression_level: 3,
encrypt: false,
encryption_key_env: None,
retention_days: 30,
max_total_size: "1KB".to_string(),
cleanup_interval_minutes: 60,
batch_size: 100,
flush_interval_ms: 100,
masking_enabled: true,
};
let result = FileSink::perform_cleanup(&config, &dir.path().join("test.log"));
assert!(result.is_ok());
let remaining = std::fs::read_dir(dir.path()).unwrap().count();
assert!(
remaining < 5,
"expected some files removed, got {}",
remaining
);
}
#[test]
fn test_perform_cleanup_empty_directory() {
let dir = tempdir().unwrap();
let log_path = dir.path().join("test.log");
let config = FileSinkConfig {
enabled: true,
path: log_path,
max_total_size: "1GB".to_string(),
..Default::default()
};
let result = FileSink::perform_cleanup(&config, &dir.path().join("test.log"));
assert!(result.is_ok());
}
#[test]
fn test_perform_cleanup_nonexistent_parent_returns_error() {
let dir = tempdir().unwrap();
let nonexistent_parent = dir.path().join("does_not_exist");
let log_path = nonexistent_parent.join("test.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
max_total_size: "1GB".to_string(),
..Default::default()
};
let result = FileSink::perform_cleanup(&config, &log_path);
assert!(result.is_err());
}
#[test]
fn test_file_sink_clone_produces_independent_instance() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
max_size: "1MB".to_string(),
rotation_time: "daily".to_string(),
..Default::default()
};
let sink = create_test_file_sink(config);
let cloned = sink.clone();
assert!(cloned.inner.read().current_file.is_none());
assert_eq!(cloned.inner.read().current_size, 0);
assert_eq!(cloned.inner.read().sequence, 0);
assert_eq!(sink.config.path, cloned.config.path);
assert_eq!(sink.config.max_size, cloned.config.max_size);
}
#[test]
fn test_file_sink_debug_format_contains_key_fields() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
let sink = create_test_file_sink(config);
let debug_str = format!("{:?}", sink);
assert!(debug_str.contains("FileSink"));
assert!(debug_str.contains("path"));
assert!(debug_str.contains("current_size"));
}
#[test]
fn test_file_sink_is_healthy_false_without_file() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
let sink = create_test_file_sink(config);
assert!(!sink.is_healthy());
}
#[test]
fn test_file_sink_is_healthy_true_with_file() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
let sink = create_test_file_sink(config);
{
let mut inner = sink.inner.write();
sink.open_file_inner(&mut inner).unwrap();
}
assert!(sink.is_healthy());
}
#[tokio::test]
async fn test_file_sink_flush_writes_buffered_records() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("test.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
..Default::default()
};
let sink = create_test_file_sink(config);
{
let mut inner = sink.inner.write();
sink.open_file_inner(&mut inner).unwrap();
inner.batch_buffer.push(create_test_record("Flush test"));
}
let result = sink.flush().await;
assert!(result.is_ok());
let content = std::fs::read_to_string(&log_path).unwrap();
assert!(content.contains("Flush test"));
}
#[tokio::test]
async fn test_file_sink_flush_without_file_succeeds() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.flush().await;
assert!(result.is_ok());
}
#[tokio::test]
async fn test_file_sink_shutdown_flushes_remaining_data() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("test.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
..Default::default()
};
let sink = create_test_file_sink(config);
{
let mut inner = sink.inner.write();
sink.open_file_inner(&mut inner).unwrap();
inner.batch_buffer.push(create_test_record("Shutdown test"));
}
let result = sink.shutdown().await;
assert!(result.is_ok());
let content = std::fs::read_to_string(&log_path).unwrap();
assert!(content.contains("Shutdown test"));
}
#[tokio::test]
async fn test_write_multiple_records_all_persisted() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("test.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
batch_size: 2, flush_interval_ms: 1000,
..Default::default()
};
let sink = FileSink::new(config).unwrap();
for i in 0..5 {
let record = create_test_record(&format!("Message {}", i));
sink.write(&record).await.unwrap();
}
sink.flush().await.unwrap();
let content = std::fs::read_to_string(&log_path).unwrap();
for i in 0..5 {
assert!(content.contains(&format!("Message {}", i)));
}
sink.shutdown().await.unwrap();
}
#[tokio::test]
async fn test_write_with_masking_disabled_preserves_sensitive_value() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("test.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
masking_enabled: false,
batch_size: 1,
..Default::default()
};
let sink = FileSink::new(config).unwrap();
let record = LogRecord {
timestamp: Utc::now(),
level: "INFO".to_string(),
target: "test".to_string(),
message: "password=secret1234567890".to_string(),
fields: HashMap::new(),
file: None,
line: None,
thread_id: "t1".to_string(),
};
sink.write(&record).await.unwrap();
sink.flush().await.unwrap();
let content = std::fs::read_to_string(&log_path).unwrap();
assert!(
content.contains("secret1234567890"),
"masking disabled should preserve original value"
);
sink.shutdown().await.unwrap();
}
#[tokio::test]
async fn test_write_with_masking_enabled_redacts_sensitive_value() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("test.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
masking_enabled: true,
batch_size: 1,
..Default::default()
};
let sink = FileSink::new(config).unwrap();
let record = LogRecord {
timestamp: Utc::now(),
level: "INFO".to_string(),
target: "test".to_string(),
message: "password=secret1234567890".to_string(),
fields: HashMap::new(),
file: None,
line: None,
thread_id: "t1".to_string(),
};
sink.write(&record).await.unwrap();
sink.flush().await.unwrap();
let content = std::fs::read_to_string(&log_path).unwrap();
assert!(
!content.contains("secret1234567890"),
"masking enabled should redact sensitive value"
);
assert!(
content.contains("***REDACTED***"),
"masked output should contain REDACTED marker"
);
sink.shutdown().await.unwrap();
}
#[tokio::test]
async fn test_write_appends_to_existing_file() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("test.log");
std::fs::write(&log_path, "pre-existing line\n").unwrap();
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
batch_size: 1,
..Default::default()
};
let sink = FileSink::new(config).unwrap();
sink.write(&create_test_record("Appended message"))
.await
.unwrap();
sink.flush().await.unwrap();
let content = std::fs::read_to_string(&log_path).unwrap();
assert!(content.starts_with("pre-existing line"));
assert!(content.contains("Appended message"));
sink.shutdown().await.unwrap();
}
#[tokio::test]
async fn test_file_sink_new_with_weekly_rotation() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("weekly.log"),
rotation_time: "weekly".to_string(),
..Default::default()
};
let sink = FileSink::new(config).unwrap();
assert_eq!(sink.rotation_interval, StdDuration::from_secs(604800));
sink.shutdown().await.unwrap();
}
#[tokio::test]
async fn test_file_sink_new_with_monthly_rotation() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("monthly.log"),
rotation_time: "monthly".to_string(),
..Default::default()
};
let sink = FileSink::new(config).unwrap();
assert_eq!(sink.rotation_interval, StdDuration::from_secs(2592000));
sink.shutdown().await.unwrap();
}
#[tokio::test]
async fn test_file_sink_new_with_unknown_rotation_falls_back_to_daily() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("unknown.log"),
rotation_time: "unknown_interval".to_string(),
..Default::default()
};
let sink = FileSink::new(config).unwrap();
assert_eq!(sink.rotation_interval, StdDuration::from_secs(86400));
sink.shutdown().await.unwrap();
}
#[tokio::test]
async fn test_file_sink_new_with_hourly_rotation() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("hourly.log"),
rotation_time: "hourly".to_string(),
..Default::default()
};
let sink = FileSink::new(config).unwrap();
assert_eq!(sink.rotation_interval, StdDuration::from_secs(3600));
sink.shutdown().await.unwrap();
}
#[test]
#[serial]
fn test_get_encryption_key_invalid_base64() {
let config = FileSinkConfig {
enabled: true,
path: PathBuf::from("test.log"),
encryption_key_env: Some("TEST_INVALID_B64".to_string()),
..Default::default()
};
std::env::set_var("TEST_INVALID_B64", "this_is_not_valid_base64!!!@#$");
let sink = create_test_file_sink(config);
let result = sink.get_encryption_key();
assert!(result.is_err());
assert!(result
.unwrap_err()
.to_string()
.contains("Invalid base64 encoding"));
std::env::remove_var("TEST_INVALID_B64");
}
#[test]
#[serial]
fn test_get_encryption_key_wrong_byte_length() {
let config = FileSinkConfig {
enabled: true,
path: PathBuf::from("test.log"),
encryption_key_env: Some("TEST_WRONG_LEN".to_string()),
..Default::default()
};
std::env::set_var("TEST_WRONG_LEN", "YWJjZGVmZ2hpamtsbW5v");
let sink = create_test_file_sink(config);
let result = sink.get_encryption_key();
assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("32 bytes"));
std::env::remove_var("TEST_WRONG_LEN");
}
#[test]
fn test_get_disk_space_info_nonexistent_path() {
let config = FileSinkConfig {
enabled: true,
path: PathBuf::from("/nonexistent_root_path_xyz/log.log"),
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.get_disk_space_info();
assert!(result.is_err());
}
#[test]
fn test_perform_cleanup_with_empty_directory() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("app.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
retention_days: 7,
max_total_size: "1GB".to_string(),
..Default::default()
};
let result = FileSink::perform_cleanup(&config, &log_path);
assert!(result.is_ok());
}
#[test]
fn test_perform_cleanup_removes_expired_files() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("app.log");
let old_file = temp_dir.path().join("app_20250101_000000.log");
std::fs::write(&old_file, "old log content").unwrap();
let old_time =
std::time::SystemTime::now() - std::time::Duration::from_secs(30 * 24 * 60 * 60);
let _ = filetime::set_file_mtime(&old_file, filetime::FileTime::from_system_time(old_time));
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
retention_days: 7, max_total_size: "1GB".to_string(),
keep_files: 0,
..Default::default()
};
let result = FileSink::perform_cleanup(&config, &log_path);
assert!(result.is_ok());
assert!(!old_file.exists(), "expired file should be removed");
}
#[test]
fn test_compress_file_basic() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("to_compress.log");
std::fs::write(&log_path, "some log content to compress\n").unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("active.log"),
compress: false, encrypt: false,
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.compress_file(&log_path);
assert!(result.is_ok(), "compress_file should succeed");
let compressed_path = result.unwrap();
assert!(compressed_path.exists(), "compressed file should exist");
assert!(compressed_path.extension().is_some_and(|e| e == "zst"));
assert!(
!log_path.exists(),
"original file should be removed after compression"
);
}
#[test]
fn test_compress_file_nonexistent_input() {
let temp_dir = tempdir().unwrap();
let nonexistent = temp_dir.path().join("does_not_exist.log");
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("active.log"),
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.compress_file(&nonexistent);
assert!(
result.is_err(),
"compress_file should fail for nonexistent input"
);
}
#[test]
#[serial]
fn test_encrypt_file_basic() {
let temp_dir = tempdir().unwrap();
let input_path = temp_dir.path().join("to_encrypt.log");
let output_path = temp_dir.path().join("encrypted.log.enc");
std::fs::write(&input_path, "secret log content\n").unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("active.log"),
encryption_key_env: Some("TEST_ENCRYPT_KEY".to_string()),
..Default::default()
};
std::env::set_var(
"TEST_ENCRYPT_KEY",
"YWJjZGVmZ2hpamtsbW5vcHFyc3R1dnd4eXoxMjM0NTY=",
);
let sink = create_test_file_sink(config);
let result = sink.encrypt_file(&input_path, &output_path);
assert!(result.is_ok(), "encrypt_file should succeed");
assert!(output_path.exists(), "encrypted file should be created");
let encrypted_size = std::fs::metadata(&output_path).unwrap().len();
assert!(
encrypted_size > 12,
"encrypted file should contain nonce + ciphertext"
);
std::env::remove_var("TEST_ENCRYPT_KEY");
}
#[test]
#[serial]
fn test_encrypt_file_missing_key_env() {
let temp_dir = tempdir().unwrap();
let input_path = temp_dir.path().join("to_encrypt.log");
let output_path = temp_dir.path().join("encrypted.log.enc");
std::fs::write(&input_path, "content\n").unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("active.log"),
encryption_key_env: Some("MISSING_ENCRYPT_KEY_ENV_VAR".to_string()),
..Default::default()
};
std::env::remove_var("MISSING_ENCRYPT_KEY_ENV_VAR");
let sink = create_test_file_sink(config);
let result = sink.encrypt_file(&input_path, &output_path);
assert!(result.is_err(), "encrypt_file should fail without key");
assert!(result
.unwrap_err()
.to_string()
.contains("Encryption key not found"));
}
#[test]
#[serial]
fn test_encrypt_file_nonexistent_input() {
let temp_dir = tempdir().unwrap();
let input_path = temp_dir.path().join("does_not_exist.log");
let output_path = temp_dir.path().join("out.log.enc");
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("active.log"),
encryption_key_env: Some("TEST_ENCRYPT_KEY_2".to_string()),
..Default::default()
};
std::env::set_var(
"TEST_ENCRYPT_KEY_2",
"YWJjZGVmZ2hpamtsbW5vcHFyc3R1dnd4eXoxMjM0NTY=",
);
let sink = create_test_file_sink(config);
let result = sink.encrypt_file(&input_path, &output_path);
assert!(
result.is_err(),
"encrypt_file should fail for nonexistent input"
);
std::env::remove_var("TEST_ENCRYPT_KEY_2");
}
#[test]
#[serial]
fn test_compress_file_with_encryption() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("to_compress_enc.log");
std::fs::write(&log_path, "content to compress and encrypt\n").unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("active.log"),
encrypt: true,
encryption_key_env: Some("TEST_COMPRESS_ENC_KEY".to_string()),
..Default::default()
};
std::env::set_var(
"TEST_COMPRESS_ENC_KEY",
"YWJjZGVmZ2hpamtsbW5vcHFyc3R1dnd4eXoxMjM0NTY=",
);
let sink = create_test_file_sink(config);
let result = sink.compress_file(&log_path);
assert!(
result.is_ok(),
"compress_file with encryption should succeed"
);
let encrypted_path = result.unwrap();
assert!(
encrypted_path.exists(),
"encrypted compressed file should exist"
);
assert!(encrypted_path.extension().is_some_and(|e| e == "enc"));
std::env::remove_var("TEST_COMPRESS_ENC_KEY");
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn test_rotate_inner_basic() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("rotate.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
compress: false,
encrypt: false,
..Default::default()
};
std::fs::write(&log_path, "original content\n").unwrap();
let sink = FileSink::new(config).unwrap();
let mut inner = sink.inner.write();
let result = sink.rotate_inner(&mut inner);
assert!(result.is_ok(), "rotate_inner should succeed");
drop(inner);
sink.shutdown().await.unwrap();
let entries: Vec<_> = std::fs::read_dir(temp_dir.path()).unwrap().collect();
assert!(
!entries.is_empty(),
"rotated file should exist in directory"
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn test_rotate_inner_with_compression() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("rotate_compress.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
compress: true,
encrypt: false,
compression_level: 3,
..Default::default()
};
std::fs::write(&log_path, "content to be rotated and compressed\n").unwrap();
let sink = FileSink::new(config).unwrap();
let mut inner = sink.inner.write();
let result = sink.rotate_inner(&mut inner);
assert!(
result.is_ok(),
"rotate_inner with compression should succeed"
);
drop(inner);
std::thread::sleep(std::time::Duration::from_millis(500));
sink.shutdown().await.unwrap();
let has_zst = std::fs::read_dir(temp_dir.path())
.unwrap()
.any(|e| e.is_ok_and(|entry| entry.path().extension().is_some_and(|ext| ext == "zst")));
assert!(has_zst, "compressed rotated file (.zst) should exist");
}
#[tokio::test]
#[serial]
#[allow(clippy::await_holding_lock)]
async fn test_rotate_inner_with_encryption_only_branch() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("rotate_encrypt.log");
let (_key_bytes, key_b64) = make_test_key();
std::env::set_var("TEST_ROTATE_ENC_KEY", &key_b64);
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
compress: false, encrypt: true, encryption_key_env: Some("TEST_ROTATE_ENC_KEY".to_string()),
..Default::default()
};
std::fs::write(&log_path, "content to be rotated and encrypted\n").unwrap();
let sink = FileSink::new(config).unwrap();
let mut inner = sink.inner.write();
let result = sink.rotate_inner(&mut inner);
assert!(
result.is_ok(),
"rotate_inner with encryption-only should succeed"
);
drop(inner);
std::thread::sleep(std::time::Duration::from_millis(500));
sink.shutdown().await.unwrap();
let has_enc = std::fs::read_dir(temp_dir.path())
.unwrap()
.any(|e| e.is_ok_and(|entry| entry.path().extension().is_some_and(|ext| ext == "enc")));
assert!(
has_enc,
"encrypted rotated file (.enc) should exist in encrypt-only mode"
);
std::env::remove_var("TEST_ROTATE_ENC_KEY");
}
#[test]
#[serial]
fn test_compress_file_with_encryption_failure_keeps_compressed() {
let temp_dir = tempdir().unwrap();
let original_path = temp_dir.path().join("to_compress_fail.log");
std::fs::write(&original_path, "content for failed encryption\n").unwrap();
let invalid_key = base64::engine::general_purpose::STANDARD.encode(b"1234567890123456");
std::env::set_var("TEST_COMPRESS_ENC_FAIL_KEY", &invalid_key);
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("active.log"),
compress: true,
compression_level: 3,
encrypt: true,
encryption_key_env: Some("TEST_COMPRESS_ENC_FAIL_KEY".to_string()),
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.compress_file(&original_path);
assert!(
result.is_err(),
"compress_file should fail when encryption key is invalid"
);
let unencrypted_file = std::fs::read_dir(temp_dir.path())
.unwrap()
.filter_map(|e| e.ok())
.map(|e| e.path())
.find(|p| {
p.to_string_lossy()
.ends_with(".unencrypted")
})
.expect("compressed file should be preserved with .unencrypted suffix when encryption fails");
let compressed_file = std::fs::File::open(&unencrypted_file).unwrap();
let decoder_result = zstd::stream::Decoder::new(compressed_file);
assert!(
decoder_result.is_ok(),
"preserved file should be valid zst compressed data"
);
std::env::remove_var("TEST_COMPRESS_ENC_FAIL_KEY");
}
#[test]
fn test_open_file_inner_create_dir_failure_returns_error() {
let temp_dir = tempdir().unwrap();
let blocking_file = temp_dir.path().join("blocking_file");
std::fs::write(&blocking_file, "block").unwrap();
let impossible_path = blocking_file.join("sub").join("log.log");
let config = FileSinkConfig {
enabled: true,
path: impossible_path,
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
let result = sink.open_file_inner(&mut inner);
assert!(
result.is_err(),
"open_file_inner should fail when parent directory cannot be created"
);
assert!(inner.current_file.is_none());
}
#[test]
fn test_file_sink_new_open_file_failure_returns_error() {
let temp_dir = tempdir().unwrap();
let blocking_file = temp_dir.path().join("block_new");
std::fs::write(&blocking_file, "block").unwrap();
let impossible_log_path = blocking_file.join("nested").join("log.log");
let config = FileSinkConfig {
enabled: true,
path: impossible_log_path,
..Default::default()
};
let result = FileSink::new(config);
assert!(
result.is_err(),
"FileSink::new should return error when open_file_inner fails"
);
}
#[test]
fn test_check_rotation_inner_by_time_triggers_rotation() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
max_size: "1MB".to_string(), rotation_time: "daily".to_string(),
compress: false,
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
sink.open_file_inner(&mut inner).unwrap();
std::fs::write(sink.config.path.clone(), "x").unwrap();
inner.next_rotation_time = Some(Utc::now() - chrono::Duration::hours(1));
let result = sink.check_rotation_inner(&mut inner);
assert!(result.is_ok());
assert_eq!(inner.sequence, 1, "rotation should be triggered by time");
}
#[test]
fn test_should_rotate_by_time_inner_weekly_date_change() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
rotation_time: "weekly".to_string(),
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
let last_week = Utc::now().date_naive().num_days_from_ce() - 7;
inner.last_rotation_date = Some(last_week);
inner.next_rotation_time = Some(Utc::now() + chrono::Duration::days(1));
let result = sink.should_rotate_by_time_inner(&inner);
assert!(
result,
"weekly rotation should trigger when date changed since last rotation"
);
}
#[test]
fn test_flush_batch_inner_write_error_records_failure_and_reopens() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("test.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
sink.open_file_inner(&mut inner).unwrap();
let _ = inner.current_file.take();
inner.batch_buffer.push(create_test_record("Will fail"));
let initial_failures = inner.circuit_breaker.failure_count();
let result = sink.flush_batch_inner(&mut inner);
assert!(
result.is_ok(),
"flush should succeed even without file handle"
);
assert_eq!(
inner.circuit_breaker.failure_count(),
initial_failures,
"no failure should be recorded when there is no file handle"
);
}
struct MockFallbackSink {
write_count: Arc<parking_lot::Mutex<usize>>,
}
struct MockFallbackSinkWithShutdownFlag {
shutdown_called: Arc<parking_lot::Mutex<bool>>,
}
#[async_trait::async_trait]
impl LogSink for MockFallbackSinkWithShutdownFlag {
async fn write(&self, _record: &LogRecord) -> Result<(), InklogError> {
Ok(())
}
async fn flush(&self) -> Result<(), InklogError> {
Ok(())
}
fn is_healthy(&self) -> bool {
true
}
async fn shutdown(&self) -> Result<(), InklogError> {
*self.shutdown_called.lock() = true;
Ok(())
}
}
#[async_trait::async_trait]
impl LogSink for MockFallbackSink {
async fn write(&self, _record: &LogRecord) -> Result<(), InklogError> {
*self.write_count.lock() += 1;
Ok(())
}
async fn flush(&self) -> Result<(), InklogError> {
Ok(())
}
fn is_healthy(&self) -> bool {
true
}
async fn shutdown(&self) -> Result<(), InklogError> {
Ok(())
}
}
#[tokio::test]
async fn test_write_with_open_circuit_breaker_uses_fallback_sink() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("test.log"),
..Default::default()
};
let sink = create_test_file_sink(config);
let write_count = Arc::new(parking_lot::Mutex::new(0usize));
let mock_sink = MockFallbackSink {
write_count: write_count.clone(),
};
{
let mut inner = sink.inner.write();
inner.fallback_sink = Some(Arc::new(mock_sink));
for _ in 0..5 {
inner.circuit_breaker.record_failure();
}
assert_eq!(inner.circuit_breaker.state(), CircuitState::Open);
}
let record = create_test_record("Fallback test");
let result = sink.write(&record).await;
assert!(
result.is_ok(),
"write should not error when circuit is open"
);
assert_eq!(
*write_count.lock(),
1,
"fallback sink should be called once when circuit breaker is open"
);
}
#[tokio::test]
async fn test_write_with_rotation_failure_uses_fallback_sink() {
let temp_dir = tempdir().unwrap();
let blocking_file = temp_dir.path().join("block_rotate");
std::fs::write(&blocking_file, "block").unwrap();
let impossible_log_path = blocking_file.join("inner.log");
let config = FileSinkConfig {
enabled: true,
path: impossible_log_path,
max_size: "1".to_string(), compress: false,
..Default::default()
};
let sink = create_test_file_sink(config);
let write_count = Arc::new(parking_lot::Mutex::new(0usize));
let mock_sink = MockFallbackSink {
write_count: write_count.clone(),
};
{
let mut inner = sink.inner.write();
inner.fallback_sink = Some(Arc::new(mock_sink));
}
let record = create_test_record("Rotation failure test");
let result = sink.write(&record).await;
assert!(result.is_ok(), "write should not error when rotation fails");
assert!(
*write_count.lock() >= 1,
"fallback sink should be called when rotation fails"
);
}
#[tokio::test]
async fn test_shutdown_with_active_timers_completes_successfully() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("shutdown_test.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
..Default::default()
};
let sink = FileSink::new(config).unwrap();
for i in 0..3 {
let record = create_test_record(&format!("Pre-shutdown message {}", i));
sink.write(&record).await.unwrap();
}
let result = sink.shutdown().await;
assert!(result.is_ok(), "shutdown should complete successfully");
let content = std::fs::read_to_string(&log_path).unwrap();
for i in 0..3 {
assert!(
content.contains(&format!("Pre-shutdown message {}", i)),
"all buffered records should be flushed before shutdown completes"
);
}
}
#[tokio::test]
async fn test_drop_does_not_panic_with_active_timers() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("drop_test.log");
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
..Default::default()
};
let sink = FileSink::new(config).unwrap();
sink.write(&create_test_record("Drop test message"))
.await
.unwrap();
drop(sink);
assert!(log_path.exists(), "log file should exist after drop");
}
#[test]
fn test_perform_cleanup_with_keep_files_boundary() {
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("keep_test.log");
let old_time =
std::time::SystemTime::now() - std::time::Duration::from_secs(30 * 24 * 60 * 60);
for i in 0..4 {
let p = temp_dir.path().join(format!("keep_{}.log", i));
std::fs::write(&p, "old content").unwrap();
let _ = filetime::set_file_mtime(&p, filetime::FileTime::from_system_time(old_time));
}
let config = FileSinkConfig {
enabled: true,
path: log_path,
retention_days: 7, keep_files: 2, max_total_size: "1GB".to_string(), ..Default::default()
};
let result = FileSink::perform_cleanup(&config, &temp_dir.path().join("keep_test.log"));
assert!(result.is_ok());
let remaining: Vec<_> = std::fs::read_dir(temp_dir.path())
.unwrap()
.filter_map(|e| e.ok())
.filter(|e| e.path().extension().is_some_and(|ext| ext == "log"))
.collect();
assert!(
remaining.len() >= 2,
"keep_files should preserve at least 2 files, got {}",
remaining.len()
);
}
#[test]
fn test_parse_size_large_values() {
assert_eq!(
FileSink::parse_size("1024TB"),
Some(1024 * 1024 * 1024 * 1024 * 1024)
);
assert_eq!(FileSink::parse_size("1KB"), Some(1024));
assert_eq!(FileSink::parse_size("1MB"), Some(1024 * 1024));
assert_eq!(FileSink::parse_size("1GB"), Some(1024 * 1024 * 1024));
assert_eq!(FileSink::parse_size("1TB"), Some(1024_u64.pow(4)));
}
#[test]
fn test_validate_key_entropy_single_byte_repeated() {
let weak_key = [0x42; 32];
let result = FileSink::validate_key_entropy(&weak_key);
assert!(
result.is_err(),
"single-byte repeated key should be rejected"
);
}
#[test]
fn test_validate_key_entropy_two_byte_pattern() {
let mut pattern_key = [0u8; 32];
for (i, byte) in pattern_key.iter_mut().enumerate() {
*byte = if i % 2 == 0 { 0xAA } else { 0x55 };
}
let result = FileSink::validate_key_entropy(&pattern_key);
assert!(
result.is_err(),
"two-byte pattern key should be rejected (entropy < 4.0)"
);
}
#[test]
fn test_validate_key_entropy_four_byte_pattern() {
let pattern = [0x11, 0x22, 0x33, 0x44];
let mut pattern_key = [0u8; 32];
for (i, byte) in pattern_key.iter_mut().enumerate() {
*byte = pattern[i % 4];
}
let result = FileSink::validate_key_entropy(&pattern_key);
assert!(
result.is_err(),
"four-byte pattern key should be rejected (entropy = 2.0 < 4.0)"
);
}
#[test]
fn test_open_file_inner_fails_when_path_is_directory() {
let temp_dir = tempdir().unwrap();
let dir_as_path = temp_dir.path().to_path_buf();
let config = FileSinkConfig {
enabled: true,
path: dir_as_path,
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
let result = sink.open_file_inner(&mut inner);
assert!(
result.is_err(),
"open_file_inner should fail when path is an existing directory"
);
assert!(inner.current_file.is_none());
}
#[test]
fn test_compress_file_fails_when_output_dir_readonly() {
use std::os::unix::fs::PermissionsExt;
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("readonly_test.log");
std::fs::write(&log_path, "test data").unwrap();
let original_perms = std::fs::metadata(temp_dir.path()).unwrap().permissions();
let mut readonly_perms = original_perms.clone();
readonly_perms.set_mode(0o555); std::fs::set_permissions(temp_dir.path(), readonly_perms).unwrap();
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
..Default::default()
};
let sink = create_test_file_sink(config);
let result = sink.compress_file(&log_path);
std::fs::set_permissions(temp_dir.path(), original_perms).unwrap();
match result {
Err(InklogError::IoError(_)) => { }
Ok(_) => {
let _ = std::fs::remove_file(log_path.with_extension("zst"));
}
other => panic!("expected IoError or Ok, got: {:?}", other),
}
}
#[test]
fn test_encrypt_file_fails_when_output_dir_readonly() {
use std::os::unix::fs::PermissionsExt;
let temp_dir = tempdir().unwrap();
let input_path = temp_dir.path().join("encrypt_input.bin");
std::fs::write(&input_path, b"plaintext data").unwrap();
let (_key_bytes, key_b64) = make_test_key();
let enc_key_env = "TEST_ENCRYPT_READONLY_KEY";
std::env::set_var(enc_key_env, &key_b64);
let original_perms = std::fs::metadata(temp_dir.path()).unwrap().permissions();
let mut readonly_perms = original_perms.clone();
readonly_perms.set_mode(0o555);
std::fs::set_permissions(temp_dir.path(), readonly_perms).unwrap();
let config = FileSinkConfig {
enabled: true,
path: input_path.clone(),
encrypt: true,
encryption_key_env: Some(enc_key_env.to_string()),
..Default::default()
};
let sink = create_test_file_sink(config);
let output_path = temp_dir.path().join("nonexistent_encrypted.enc");
let result = sink.encrypt_file(&input_path, &output_path);
std::fs::set_permissions(temp_dir.path(), original_perms).unwrap();
std::env::remove_var(enc_key_env);
match result {
Err(InklogError::IoError(_)) => { }
Ok(_) => {
let _ = std::fs::remove_file(&output_path);
}
other => panic!("expected IoError or Ok, got: {:?}", other),
}
}
#[test]
fn test_rotate_inner_rename_failure_returns_error_when_copy_also_fails() {
use std::os::unix::fs::PermissionsExt;
let temp_dir = tempdir().unwrap();
let log_path = temp_dir.path().join("rotate_rename_fail.log");
std::fs::write(&log_path, "rotation test data").unwrap();
let original_perms = std::fs::metadata(temp_dir.path()).unwrap().permissions();
let mut readonly_perms = original_perms.clone();
readonly_perms.set_mode(0o555);
std::fs::set_permissions(temp_dir.path(), readonly_perms).unwrap();
let config = FileSinkConfig {
enabled: true,
path: log_path.clone(),
..Default::default()
};
let sink = create_test_file_sink(config);
let mut inner = sink.inner.write();
let result = sink.rotate_inner(&mut inner);
std::fs::set_permissions(temp_dir.path(), original_perms).unwrap();
match result {
Err(InklogError::IoError(_)) => { }
Ok(_) => {
let _ = std::fs::remove_file(log_path);
}
other => panic!("expected IoError or Ok, got: {:?}", other),
}
}
#[tokio::test]
async fn test_shutdown_calls_fallback_sink_shutdown() {
let temp_dir = tempdir().unwrap();
let config = FileSinkConfig {
enabled: true,
path: temp_dir.path().join("shutdown_fallback.log"),
..Default::default()
};
let sink = create_test_file_sink(config);
let shutdown_called = Arc::new(parking_lot::Mutex::new(false));
let mock_sink = MockFallbackSinkWithShutdownFlag {
shutdown_called: shutdown_called.clone(),
};
{
let mut inner = sink.inner.write();
inner.fallback_sink = Some(Arc::new(mock_sink));
}
let result = sink.shutdown().await;
assert!(result.is_ok(), "shutdown should succeed");
assert!(
*shutdown_called.lock(),
"fallback sink shutdown should be called"
);
}
}