use anyhow::Result;
use chrono::{DateTime, Duration, Utc};
use serde::{Deserialize, Serialize};
use std::fs;
use std::io::{BufRead, BufReader, Write};
use std::path::{Path, PathBuf};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RetentionPolicy {
pub max_age_days: Option<u32>,
pub max_events: Option<usize>,
pub max_file_size_bytes: Option<u64>,
pub archive_old_events: bool,
pub archive_path: Option<PathBuf>,
pub compress_archives: bool,
}
impl Default for RetentionPolicy {
fn default() -> Self {
let archive_path = if let Ok(global_base) = crate::storage::get_default_storage_dir() {
Some(global_base.join("events").join("archive"))
} else {
Some(PathBuf::from(".prodigy/events/archive"))
};
Self {
max_age_days: Some(30), max_events: Some(100000), max_file_size_bytes: Some(100 * 1024 * 1024), archive_old_events: true,
archive_path,
compress_archives: true,
}
}
}
pub struct RetentionManager {
policy: RetentionPolicy,
events_path: PathBuf,
}
impl RetentionManager {
pub fn new(policy: RetentionPolicy, events_path: PathBuf) -> Self {
Self {
policy,
events_path,
}
}
pub fn with_default_policy(events_path: PathBuf) -> Self {
Self::new(RetentionPolicy::default(), events_path)
}
pub async fn with_global_storage(repo_path: &Path, job_id: &str) -> Result<Self> {
use crate::storage::{extract_repo_name, GlobalStorage};
let storage = GlobalStorage::new()?;
let repo_name = extract_repo_name(repo_path)?;
let events_path = storage.get_events_dir(&repo_name, job_id).await?;
Ok(Self::with_default_policy(events_path))
}
pub fn from_config_file(config_path: &Path, events_path: PathBuf) -> Result<Self> {
let config_content = fs::read_to_string(config_path)?;
let policy: RetentionPolicy = serde_yaml::from_str(&config_content)?;
Ok(Self::new(policy, events_path))
}
pub async fn analyze_retention(&self) -> Result<RetentionAnalysis> {
let mut analysis = RetentionAnalysis {
file_path: self.events_path.clone(),
..Default::default()
};
if !self.events_path.exists() {
analysis.warnings.push("File does not exist".to_string());
return Ok(analysis);
}
let metadata = fs::metadata(&self.events_path)?;
analysis.original_size_bytes = metadata.len();
let cutoff_time = self.calculate_cutoff_time();
let file = fs::File::open(&self.events_path)?;
let reader = BufReader::new(file);
let mut events_to_keep = 0usize;
let mut events_to_remove = 0usize;
let mut bytes_retained = 0u64;
let mut last_progress_report = 0usize;
const PROGRESS_REPORT_INTERVAL: usize = 10000;
for line in reader.lines() {
let line = line?;
if line.trim().is_empty() {
continue;
}
analysis.events_total += 1;
if analysis.events_total >= last_progress_report + PROGRESS_REPORT_INTERVAL {
eprint!("\rAnalyzing events: {} processed...", analysis.events_total);
use std::io::Write;
std::io::stderr().flush().ok();
last_progress_report = analysis.events_total;
}
if let Ok(event) = serde_json::from_str::<serde_json::Value>(&line) {
if self.should_retain_event(&event, cutoff_time, events_to_keep) {
events_to_keep += 1;
bytes_retained += line.len() as u64 + 1; } else {
events_to_remove += 1;
}
}
}
if last_progress_report > 0 {
eprint!("\r{}\r", " ".repeat(50));
std::io::stderr().flush().ok();
}
analysis.events_retained = events_to_keep;
analysis.events_to_remove = events_to_remove;
if self.policy.archive_old_events && events_to_remove > 0 {
analysis.events_to_archive = events_to_remove;
}
analysis.projected_size_bytes = bytes_retained;
analysis.space_to_save = analysis
.original_size_bytes
.saturating_sub(analysis.projected_size_bytes);
if analysis.events_to_remove > 10000 {
analysis.warnings.push(format!(
"Large number of events will be removed: {}",
analysis.events_to_remove
));
}
if analysis.space_to_save > 100 * 1024 * 1024 {
analysis.warnings.push(format!(
"Large amount of space will be freed: {:.1} MB",
analysis.space_to_save as f64 / (1024.0 * 1024.0)
));
}
if self.needs_cleanup(analysis.original_size_bytes)? && analysis.events_to_remove == 0 {
analysis.warnings.push("Cleanup triggered but no events would be removed - consider adjusting retention policy".to_string());
}
analysis.estimated_duration_secs = self.estimate_operation_duration(
analysis.original_size_bytes,
analysis.events_total,
analysis.events_to_remove,
self.policy.archive_old_events,
);
Ok(analysis)
}
pub async fn apply_retention(&self) -> Result<RetentionStats> {
let mut stats = RetentionStats::default();
if !self.events_path.exists() {
return Ok(stats);
}
let metadata = fs::metadata(&self.events_path)?;
let file_size = metadata.len();
stats.original_size_bytes = file_size;
let needs_cleanup = self.needs_cleanup(file_size)?;
if !needs_cleanup {
stats.events_retained = self.count_events()?;
stats.final_size_bytes = file_size;
return Ok(stats);
}
self.cleanup_events(&mut stats).await?;
Ok(stats)
}
fn needs_cleanup(&self, file_size: u64) -> Result<bool> {
if let Some(max_size) = self.policy.max_file_size_bytes {
if file_size > max_size {
return Ok(true);
}
}
if let Some(max_events) = self.policy.max_events {
let event_count = self.count_events()?;
if event_count > max_events {
return Ok(true);
}
}
if self.policy.max_age_days.is_some() {
return Ok(true);
}
Ok(false)
}
fn count_events(&self) -> Result<usize> {
let file = fs::File::open(&self.events_path)?;
let reader = BufReader::new(file);
let count = reader
.lines()
.map_while(Result::ok)
.filter(|l| !l.trim().is_empty())
.count();
Ok(count)
}
async fn cleanup_events(&self, stats: &mut RetentionStats) -> Result<()> {
let cutoff_time = self.calculate_cutoff_time();
let temp_file = self.events_path.with_extension("tmp");
let mut events_to_archive = Vec::new();
let mut events_to_keep = Vec::new();
let file = fs::File::open(&self.events_path)?;
let reader = BufReader::new(file);
for line in reader.lines() {
let line = line?;
if line.trim().is_empty() {
continue;
}
stats.events_processed += 1;
if let Ok(event) = serde_json::from_str::<serde_json::Value>(&line) {
if self.should_retain_event(&event, cutoff_time, stats.events_retained) {
events_to_keep.push(line);
stats.events_retained += 1;
} else {
events_to_archive.push(line);
stats.events_removed += 1;
}
}
}
if self.policy.archive_old_events && !events_to_archive.is_empty() {
self.archive_events(&events_to_archive, stats).await?;
}
let mut temp_writer = fs::File::create(&temp_file)?;
for event in events_to_keep {
writeln!(temp_writer, "{}", event)?;
}
temp_writer.sync_all()?;
fs::rename(&temp_file, &self.events_path)?;
let metadata = fs::metadata(&self.events_path)?;
stats.final_size_bytes = metadata.len();
Ok(())
}
fn calculate_cutoff_time(&self) -> Option<DateTime<Utc>> {
self.policy
.max_age_days
.map(|days| Utc::now() - Duration::days(days as i64))
}
fn should_retain_event(
&self,
event: &serde_json::Value,
cutoff_time: Option<DateTime<Utc>>,
current_retained_count: usize,
) -> bool {
if let Some(max_events) = self.policy.max_events {
if current_retained_count >= max_events {
return false;
}
}
if let Some(cutoff) = cutoff_time {
if let Some(timestamp) = extract_event_timestamp(event) {
if timestamp < cutoff {
return false;
}
}
}
true
}
async fn archive_events(&self, events: &[String], stats: &mut RetentionStats) -> Result<()> {
let archive_dir = self
.policy
.archive_path
.as_ref()
.ok_or_else(|| anyhow::anyhow!("Archive path not configured"))?;
fs::create_dir_all(archive_dir)?;
let archive_filename = format!(
"events_archive_{}.jsonl{}",
Utc::now().format("%Y%m%d_%H%M%S"),
if self.policy.compress_archives {
".gz"
} else {
""
}
);
let archive_path = archive_dir.join(archive_filename);
if self.policy.compress_archives {
self.write_compressed_archive(&archive_path, events)?;
} else {
self.write_plain_archive(&archive_path, events)?;
}
stats.events_archived = events.len();
stats.archive_path = Some(archive_path);
Ok(())
}
fn write_plain_archive(&self, path: &Path, events: &[String]) -> Result<()> {
let mut file = fs::File::create(path)?;
for event in events {
writeln!(file, "{}", event)?;
}
file.sync_all()?;
Ok(())
}
fn write_compressed_archive(&self, path: &Path, events: &[String]) -> Result<()> {
use flate2::write::GzEncoder;
use flate2::Compression;
let file = fs::File::create(path)?;
let mut encoder = GzEncoder::new(file, Compression::default());
for event in events {
writeln!(encoder, "{}", event)?;
}
encoder.finish()?;
Ok(())
}
pub fn policy(&self) -> &RetentionPolicy {
&self.policy
}
pub fn set_policy(&mut self, policy: RetentionPolicy) {
self.policy = policy;
}
pub fn save_policy_to_file(&self, config_path: &Path) -> Result<()> {
let yaml = serde_yaml::to_string(&self.policy)?;
fs::write(config_path, yaml)?;
Ok(())
}
fn estimate_operation_duration(
&self,
file_size_bytes: u64,
total_events: usize,
events_to_process: usize,
archive_enabled: bool,
) -> f64 {
const BASE_OVERHEAD_SECS: f64 = 0.5;
const BYTES_PER_SEC_READ: f64 = 50_000_000.0; const BYTES_PER_SEC_WRITE: f64 = 30_000_000.0; const EVENTS_PER_SEC_PROCESS: f64 = 10_000.0; const ARCHIVE_OVERHEAD_FACTOR: f64 = 1.5;
let read_time = (file_size_bytes as f64) / BYTES_PER_SEC_READ;
let processing_time = (total_events as f64) / EVENTS_PER_SEC_PROCESS;
let retention_ratio = 1.0 - (events_to_process as f64 / total_events.max(1) as f64);
let write_size = (file_size_bytes as f64) * retention_ratio;
let write_time = write_size / BYTES_PER_SEC_WRITE;
let archive_time = if archive_enabled && events_to_process > 0 {
let archive_size =
(file_size_bytes as f64) * (events_to_process as f64 / total_events.max(1) as f64);
(archive_size / BYTES_PER_SEC_WRITE) * ARCHIVE_OVERHEAD_FACTOR
} else {
0.0
};
BASE_OVERHEAD_SECS + read_time + processing_time + write_time + archive_time
}
}
#[derive(Debug, Default, Clone)]
pub struct RetentionStats {
pub events_processed: usize,
pub events_retained: usize,
pub events_removed: usize,
pub events_archived: usize,
pub original_size_bytes: u64,
pub final_size_bytes: u64,
pub archive_path: Option<PathBuf>,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
pub struct RetentionAnalysis {
pub file_path: PathBuf,
pub events_total: usize,
pub events_retained: usize,
pub events_to_remove: usize,
pub events_to_archive: usize,
pub original_size_bytes: u64,
pub projected_size_bytes: u64,
pub space_to_save: u64,
pub estimated_duration_secs: f64,
pub warnings: Vec<String>,
}
impl RetentionAnalysis {
pub fn display_human(&self) {
println!("Cleanup Analysis (DRY RUN)");
println!("========================");
println!("File: {}", self.file_path.display());
println!("Total events: {}", self.events_total);
println!("Events to retain: {}", self.events_retained);
println!("Events to remove: {}", self.events_to_remove);
if self.events_to_archive > 0 {
println!("Events to archive: {}", self.events_to_archive);
}
println!("Current size: {} bytes", self.original_size_bytes);
println!("Projected size: {} bytes", self.projected_size_bytes);
println!(
"Space to save: {} bytes ({:.1}%)",
self.space_to_save,
if self.original_size_bytes > 0 {
(self.space_to_save as f64 / self.original_size_bytes as f64) * 100.0
} else {
0.0
}
);
if self.estimated_duration_secs > 0.0 {
println!(
"Estimated time: {}",
format_duration(self.estimated_duration_secs)
);
}
if !self.warnings.is_empty() {
println!("\nWarnings:");
for warning in &self.warnings {
println!(" ⚠️ {}", warning);
}
}
}
}
fn format_duration(secs: f64) -> String {
if secs < 1.0 {
format!("{:.0} ms", secs * 1000.0)
} else if secs < 60.0 {
format!("{:.1} seconds", secs)
} else if secs < 3600.0 {
let mins = secs / 60.0;
format!("{:.1} minutes", mins)
} else {
let hours = secs / 3600.0;
format!("{:.1} hours", hours)
}
}
impl RetentionStats {
pub fn space_saved(&self) -> u64 {
self.original_size_bytes
.saturating_sub(self.final_size_bytes)
}
pub fn space_saved_percentage(&self) -> f64 {
if self.original_size_bytes > 0 {
(self.space_saved() as f64 / self.original_size_bytes as f64) * 100.0
} else {
0.0
}
}
pub fn display_summary(&self) {
println!("Event Retention Summary:");
println!(" Events processed: {}", self.events_processed);
println!(" Events retained: {}", self.events_retained);
println!(" Events removed: {}", self.events_removed);
if self.events_archived > 0 {
println!(" Events archived: {}", self.events_archived);
if let Some(ref path) = self.archive_path {
println!(" Archive location: {}", path.display());
}
}
println!(" Original size: {} bytes", self.original_size_bytes);
println!(" Final size: {} bytes", self.final_size_bytes);
println!(
" Space saved: {} bytes ({:.1}%)",
self.space_saved(),
self.space_saved_percentage()
);
}
}
fn extract_event_timestamp(event: &serde_json::Value) -> Option<DateTime<Utc>> {
let timestamp_str = event
.get("timestamp")
.or_else(|| event.get("time"))
.or_else(|| event.get("created_at"))
.or_else(|| {
for key in [
"JobStarted",
"JobCompleted",
"AgentStarted",
"AgentCompleted",
] {
if let Some(nested) = event.get(key) {
if let Some(ts) = nested.get("timestamp") {
return Some(ts);
}
}
}
None
})
.and_then(|v| v.as_str());
timestamp_str
.and_then(|ts| DateTime::parse_from_rfc3339(ts).ok())
.map(|dt| dt.with_timezone(&Utc))
}
pub struct RetentionTask {
manager: RetentionManager,
interval: std::time::Duration,
}
impl RetentionTask {
pub fn new(manager: RetentionManager, interval: std::time::Duration) -> Self {
Self { manager, interval }
}
pub async fn run_once(&self) -> Result<RetentionStats> {
log::info!("Running event retention cleanup...");
let stats = self.manager.apply_retention().await?;
if stats.events_removed > 0 {
log::info!(
"Retention cleanup completed: {} events removed, {:.1}% space saved",
stats.events_removed,
stats.space_saved_percentage()
);
} else {
log::debug!("Retention cleanup completed: no events removed");
}
Ok(stats)
}
pub async fn start(self) {
let mut interval = tokio::time::interval(self.interval);
loop {
interval.tick().await;
if let Err(e) = self.run_once().await {
log::error!("Retention task failed: {}", e);
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::TempDir;
#[test]
fn test_default_retention_policy() {
let policy = RetentionPolicy::default();
assert_eq!(policy.max_age_days, Some(30));
assert_eq!(policy.max_events, Some(100000));
assert_eq!(policy.max_file_size_bytes, Some(100 * 1024 * 1024));
assert!(policy.archive_old_events);
assert!(policy.compress_archives);
}
#[tokio::test]
async fn test_retention_manager_no_cleanup_needed() {
let temp_dir = TempDir::new().unwrap();
let events_file = temp_dir.path().join("events.jsonl");
let recent_timestamp = Utc::now().to_rfc3339();
let content = format!(r#"{{"timestamp":"{}","event":"test"}}"#, recent_timestamp);
std::fs::write(&events_file, content).unwrap();
let policy = RetentionPolicy {
max_age_days: Some(365), max_events: Some(10000), max_file_size_bytes: Some(100 * 1024 * 1024), archive_old_events: false,
archive_path: None,
compress_archives: false,
};
let manager = RetentionManager::new(policy, events_file);
let stats = manager.apply_retention().await.unwrap();
assert_eq!(stats.events_processed, 1);
assert_eq!(stats.events_retained, 1);
assert_eq!(stats.events_removed, 0);
}
#[test]
fn test_extract_event_timestamp() {
let event_json = r#"{
"timestamp": "2024-01-01T12:00:00Z",
"event_type": "JobStarted"
}"#;
let event: serde_json::Value = serde_json::from_str(event_json).unwrap();
let timestamp = extract_event_timestamp(&event);
assert!(timestamp.is_some());
use chrono::Datelike;
let ts = timestamp.unwrap();
assert_eq!(ts.year(), 2024);
assert_eq!(ts.month(), 1);
assert_eq!(ts.day(), 1);
}
#[test]
fn test_retention_stats_calculations() {
let stats = RetentionStats {
original_size_bytes: 1000,
final_size_bytes: 250,
..Default::default()
};
assert_eq!(stats.space_saved(), 750);
assert_eq!(stats.space_saved_percentage(), 75.0);
}
}