use std::sync::Arc;
use crate::storage::repositories::audit_logs::AuditLogRepository;
pub fn start_audit_cleanup(
retention_days: u64,
interval_secs: u64,
audit_repo: Arc<dyn AuditLogRepository>,
lease_gate: Option<Arc<crate::cluster::JobLeaseGate>>,
) -> Option<tokio::task::JoinHandle<()>> {
if retention_days == 0 {
tracing::info!("Audit log retention disabled (audit.retention_days = 0)");
return None;
}
let handle = super::start_retention_job(
"audit_cleanup",
interval_secs,
lease_gate,
move || {
let repo = audit_repo.clone();
async move { repo.delete_older_than(retention_days).await }
},
move |outcome| match outcome {
Ok(count) => {
if count > 0 {
tracing::info!(
deleted = count,
retention_days = retention_days,
"Audit log cleanup completed"
);
}
}
Err(e) => {
tracing::error!(error = %e, "Audit log cleanup failed");
}
},
);
tracing::info!(
retention_days = retention_days,
interval_secs = interval_secs,
"Audit log cleanup task started"
);
Some(handle)
}
#[cfg(test)]
mod tests {
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
use super::*;
use crate::errors::OrionError;
use crate::storage::models::AuditLogEntry;
use crate::storage::repositories::audit_logs::AuditLogFilter;
use crate::storage::repositories::cluster::{ClusterRepository, EpochRow};
use crate::storage::repositories::helpers::PaginatedResult;
#[derive(Default)]
struct MockAuditRepo {
deletes: AtomicUsize,
}
#[async_trait::async_trait]
impl AuditLogRepository for MockAuditRepo {
async fn insert(
&self,
_principal: &str,
_action: &str,
_resource_type: &str,
_resource_id: &str,
_details: Option<&str>,
) -> Result<(), OrionError> {
Ok(())
}
async fn list_paginated(
&self,
_filter: &AuditLogFilter,
) -> Result<PaginatedResult<AuditLogEntry>, OrionError> {
Ok(PaginatedResult {
data: vec![],
total: 0,
limit: 0,
offset: 0,
})
}
async fn delete_older_than(&self, _days: u64) -> Result<u64, OrionError> {
self.deletes.fetch_add(1, Ordering::SeqCst);
Ok(1)
}
}
struct LeaseHeldElsewhere;
#[async_trait::async_trait]
impl ClusterRepository for LeaseHeldElsewhere {
async fn bump_epoch(&self) -> Result<i64, OrionError> {
unreachable!("not used by audit cleanup")
}
async fn get_epoch(&self) -> Result<EpochRow, OrionError> {
unreachable!("not used by audit cleanup")
}
async fn request_breaker_reset(&self, _key: &str) -> Result<i64, OrionError> {
unreachable!("not used by audit cleanup")
}
async fn try_acquire_job_lease(
&self,
_job_name: &str,
_holder: &str,
_ttl_secs: u64,
) -> Result<bool, OrionError> {
Ok(false)
}
}
async fn advance_and_yield(duration: Duration) {
tokio::time::advance(duration).await;
for _ in 0..20 {
tokio::task::yield_now().await;
}
}
#[tokio::test(start_paused = true)]
async fn test_disabled_when_retention_is_zero() {
let repo = Arc::new(MockAuditRepo::default());
let handle = start_audit_cleanup(0, 1, repo.clone(), None);
assert!(handle.is_none(), "0 days must mean retain forever");
advance_and_yield(Duration::from_secs(3)).await;
assert_eq!(repo.deletes.load(Ordering::SeqCst), 0);
}
#[tokio::test(start_paused = true)]
async fn test_ungated_job_deletes_expired_rows() {
let repo = Arc::new(MockAuditRepo::default());
let handle = start_audit_cleanup(90, 1, repo.clone(), None).expect("job started");
advance_and_yield(Duration::from_secs(1)).await;
advance_and_yield(Duration::from_secs(1)).await;
assert!(
repo.deletes.load(Ordering::SeqCst) >= 1,
"the first interval tick must run the DELETE"
);
handle.abort();
}
fn render_job_metrics<F, Fut>(job: F) -> String
where
F: FnOnce() -> Fut,
Fut: std::future::Future<Output = ()>,
{
let recorder = metrics_exporter_prometheus::PrometheusBuilder::new().build_recorder();
let handle = recorder.handle();
::metrics::with_local_recorder(&recorder, || {
crate::metrics::set_enabled(true);
tokio::runtime::Builder::new_current_thread()
.enable_time()
.start_paused(true)
.build()
.expect("test runtime")
.block_on(job());
});
handle.render()
}
#[test]
fn test_successful_tick_stamps_the_job_health_gauge() {
let out = render_job_metrics(|| async {
let repo = Arc::new(MockAuditRepo::default());
let handle = start_audit_cleanup(90, 1, repo, None).expect("job started");
advance_and_yield(Duration::from_secs(1)).await;
advance_and_yield(Duration::from_secs(1)).await;
handle.abort();
});
assert!(
out.contains(r#"orion_job_last_success_timestamp_seconds{job="audit_cleanup"}"#),
"a successful tick must stamp the job health gauge:\n{out}"
);
}
#[test]
fn test_lease_refused_tick_does_not_stamp_the_gauge() {
let out = render_job_metrics(|| async {
let repo = Arc::new(MockAuditRepo::default());
let gate = Arc::new(crate::cluster::JobLeaseGate::new(
Arc::new(LeaseHeldElsewhere),
"node-b".to_string(),
));
let handle = start_audit_cleanup(90, 1, repo, Some(gate)).expect("job started");
advance_and_yield(Duration::from_secs(1)).await;
advance_and_yield(Duration::from_secs(1)).await;
handle.abort();
});
assert!(
!out.contains(r#"job="audit_cleanup""#),
"a tick skipped for the lease is not a success:\n{out}"
);
}
#[tokio::test(start_paused = true)]
async fn test_job_skips_tick_while_another_node_holds_the_lease() {
let repo = Arc::new(MockAuditRepo::default());
let gate = Arc::new(crate::cluster::JobLeaseGate::new(
Arc::new(LeaseHeldElsewhere),
"node-b".to_string(),
));
let handle = start_audit_cleanup(90, 1, repo.clone(), Some(gate)).expect("job started");
advance_and_yield(Duration::from_secs(1)).await;
for _ in 0..3 {
advance_and_yield(Duration::from_secs(1)).await;
}
handle.abort();
assert_eq!(
repo.deletes.load(Ordering::SeqCst),
0,
"the lease holder is another node, so this one must not delete"
);
}
}