use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use tracing::{debug, error, info, warn};
use super::{JobHandler, JobOutcome, JobQueue, JobSpec};
use crate::sqlite::db::Database;
use crate::sqlite::job::Job;
const SWEEP_KEY: &str = "all";
pub const RETENTION_JOB_KIND: &str = "job_retention_sweep";
pub const NONCE_SWEEP_KIND: &str = "nonce_sweep";
pub const AUDIT_SWEEP_KIND: &str = "audit_sweep";
pub const ADMIN_SESSION_SWEEP_KIND: &str = "admin_session_sweep";
const DAILY: Duration = Duration::from_secs(24 * 60 * 60);
#[derive(Debug, Clone, Copy)]
enum SweepTarget {
Nonces { ttl: Duration },
AuditLog { retention_days: u64 },
AdminSessions { idle_timeout: Duration },
Jobs { retention_days: u64 },
}
pub struct SweepJob {
target: SweepTarget,
interval: Duration,
database: Arc<Database>,
}
impl SweepJob {
#[must_use]
pub fn nonces(database: Arc<Database>, ttl: Duration) -> Self {
Self {
target: SweepTarget::Nonces { ttl },
interval: (ttl / 2).max(Duration::from_secs(30)),
database,
}
}
#[must_use]
pub fn audit(database: Arc<Database>, retention_days: u64) -> Self {
Self {
target: SweepTarget::AuditLog { retention_days },
interval: DAILY,
database,
}
}
#[must_use]
pub fn admin_sessions(
database: Arc<Database>,
idle_timeout: Duration,
ttl_seconds: u64,
) -> Self {
Self {
target: SweepTarget::AdminSessions { idle_timeout },
interval: Duration::from_secs((ttl_seconds / 4).max(60)),
database,
}
}
#[must_use]
pub fn jobs(database: Arc<Database>, retention_days: u64) -> Self {
Self {
target: SweepTarget::Jobs { retention_days },
interval: DAILY,
database,
}
}
#[must_use]
pub fn interval(&self) -> Duration {
self.interval
}
async fn sweep(&self) {
match self.target {
SweepTarget::Nonces { ttl } => {
match crate::sqlite::nonce::Nonce::cleanup(&self.database, ttl).await {
Ok(removed) => debug!(
event = "nonce_reaper_swept",
outcome = "success",
rows_removed = removed
),
Err(error) => {
error!(event = "nonce_reaper_failed", outcome = "failure", error = %error);
}
}
}
SweepTarget::AuditLog { retention_days } => {
let cutoff = crate::admin::ops::audit_cutoff(retention_days);
match crate::sqlite::audit::AuditEntry::cleanup(cutoff, &self.database).await {
Ok(removed) => info!(
event = "audit_reaper_swept",
outcome = "success",
rows_removed = removed,
cutoff
),
Err(error) => {
error!(event = "audit_reaper_failed", outcome = "failure", error = %error);
}
}
}
SweepTarget::AdminSessions { idle_timeout } => {
match crate::sqlite::admin_session::AdminSession::cleanup(
idle_timeout,
&self.database,
)
.await
{
Ok(removed) => debug!(
event = "admin_session_reaper_swept",
outcome = "success",
rows_removed = removed
),
Err(error) => {
error!(event = "admin_session_reaper_failed", outcome = "failure", error = %error);
}
}
}
SweepTarget::Jobs { retention_days } => {
let cutoff = crate::admin::ops::audit_cutoff(retention_days);
if let Err(error) = Job::cleanup(cutoff, &self.database).await {
warn!(event = "job_retention_sweep_failed", outcome = "failure", error = %error);
}
}
}
}
}
#[async_trait]
impl JobHandler for SweepJob {
fn kind(&self) -> &'static str {
match self.target {
SweepTarget::Nonces { .. } => NONCE_SWEEP_KIND,
SweepTarget::AuditLog { .. } => AUDIT_SWEEP_KIND,
SweepTarget::AdminSessions { .. } => ADMIN_SESSION_SWEEP_KIND,
SweepTarget::Jobs { .. } => RETENTION_JOB_KIND,
}
}
async fn run(&self, _job: &Job) -> JobOutcome {
self.sweep().await;
JobOutcome::Reschedule(self.interval)
}
async fn recover(&self, queue: &JobQueue) {
queue
.enqueue_or_log(JobSpec::now(self.kind(), SWEEP_KEY))
.await;
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::JobsConfig;
use crate::sqlite::nonce::{Nonce, now_secs};
use serde_json::json;
async fn setup() -> (Arc<Database>, JobQueue) {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
let queue = JobQueue::new(database.clone(), &JobsConfig::default());
(database, queue)
}
fn row(kind: &str) -> Job {
Job {
id: "sweep".to_string(),
kind: kind.to_string(),
dedup_key: SWEEP_KEY.to_string(),
payload: json!({}),
status: "running".to_string(),
run_at: now_secs(),
attempts: 1,
max_attempts: 5,
deadline: None,
lease_until: None,
lease_owner: None,
last_error: None,
created_at: now_secs(),
updated_at: now_secs(),
}
}
#[tokio::test]
async fn each_target_answers_its_own_kind() {
let (database, _queue) = setup().await;
let kinds = [
SweepJob::nonces(database.clone(), Duration::from_secs(300)).kind(),
SweepJob::audit(database.clone(), 7).kind(),
SweepJob::admin_sessions(database.clone(), Duration::from_secs(3600), 43200).kind(),
SweepJob::jobs(database, 7).kind(),
];
let unique: std::collections::BTreeSet<_> = kinds.iter().collect();
assert_eq!(unique.len(), kinds.len(), "{kinds:?}");
}
#[tokio::test]
async fn the_intervals_keep_their_floors() {
let (database, _queue) = setup().await;
assert_eq!(
SweepJob::nonces(database.clone(), Duration::from_secs(300)).interval(),
Duration::from_secs(150)
);
assert_eq!(
SweepJob::nonces(database.clone(), Duration::from_secs(10)).interval(),
Duration::from_secs(30)
);
assert_eq!(
SweepJob::admin_sessions(database.clone(), Duration::from_secs(3600), 43200).interval(),
Duration::from_secs(10800)
);
assert_eq!(
SweepJob::admin_sessions(database.clone(), Duration::from_secs(3600), 4).interval(),
Duration::from_secs(60)
);
assert_eq!(SweepJob::audit(database.clone(), 7).interval(), DAILY);
assert_eq!(SweepJob::jobs(database, 7).interval(), DAILY);
}
#[tokio::test]
async fn the_nonce_sweep_removes_expired_rows_and_reschedules() {
let (database, _queue) = setup().await;
let ttl = Duration::from_secs(300);
let fresh = Nonce::new();
fresh.save(&database).await.unwrap();
sqlx::query("INSERT INTO nonces (value, created_at) VALUES ('stale', ?);")
.bind(now_secs() - 3600)
.execute(&database.pool)
.await
.unwrap();
let handler = SweepJob::nonces(database.clone(), ttl);
match handler.run(&row(NONCE_SWEEP_KIND)).await {
JobOutcome::Reschedule(delay) => assert_eq!(delay, Duration::from_secs(150)),
other => panic!("{other:?}"),
}
assert!(
!Nonce::verify("stale", &database, ttl).await.unwrap(),
"the expired nonce is gone"
);
assert!(
Nonce::verify(&fresh.value, &database, ttl).await.unwrap(),
"a fresh nonce survives its own sweep"
);
}
#[tokio::test]
async fn the_audit_sweep_removes_rows_past_the_retention() {
let (database, _queue) = setup().await;
crate::sqlite::audit::AuditEntry::insert(
crate::audit::AuditRecord::new(
crate::audit::AuditEvent::CertificateIssued,
"default",
crate::audit::Actor::system(),
),
&database,
)
.await
.unwrap();
let stale = now_secs() - 30 * 24 * 60 * 60;
sqlx::query("UPDATE audit_log SET created_at = ?;")
.bind(stale)
.execute(&database.pool)
.await
.unwrap();
let handler = SweepJob::audit(database.clone(), 7);
assert!(matches!(
handler.run(&row(AUDIT_SWEEP_KIND)).await,
JobOutcome::Reschedule(_)
));
let (remaining,): (i64,) = sqlx::query_as("SELECT COUNT(*) FROM audit_log;")
.fetch_one(&database.pool)
.await
.unwrap();
assert_eq!(remaining, 0);
}
#[tokio::test]
async fn the_admin_session_sweep_removes_expired_rows() {
let (database, _queue) = setup().await;
crate::sqlite::admin_user::AdminUser::create("ops", "hash", &database)
.await
.unwrap();
let user = crate::sqlite::admin_user::AdminUser::find_by_username("ops", &database)
.await
.unwrap()
.unwrap();
crate::sqlite::admin_session::AdminSession::create(
crate::sqlite::admin_session::NewSession {
user_id: &user.id,
token_hash: "hash",
csrf_token: "csrf",
created_ip: None,
user_agent: None,
},
Duration::from_secs(3600),
&database,
)
.await
.unwrap();
sqlx::query("UPDATE admin_sessions SET expires_at = ?;")
.bind(now_secs() - 60)
.execute(&database.pool)
.await
.unwrap();
let handler = SweepJob::admin_sessions(database.clone(), Duration::from_secs(3600), 43200);
assert!(matches!(
handler.run(&row(ADMIN_SESSION_SWEEP_KIND)).await,
JobOutcome::Reschedule(_)
));
let (remaining,): (i64,) = sqlx::query_as("SELECT COUNT(*) FROM admin_sessions;")
.fetch_one(&database.pool)
.await
.unwrap();
assert_eq!(remaining, 0);
}
#[tokio::test]
async fn the_queue_sweep_removes_settled_rows() {
let (database, _queue) = setup().await;
let stale = now_secs() - 30 * 24 * 60 * 60;
sqlx::query(
"INSERT INTO jobs (id, kind, dedup_key, payload, status, run_at, attempts, \
max_attempts, created_at, updated_at) \
VALUES ('old', 'test', 'k', '{}', 'failed', ?, 1, 1, ?, ?);",
)
.bind(stale)
.bind(stale)
.bind(stale)
.execute(&database.pool)
.await
.unwrap();
let handler = SweepJob::jobs(database.clone(), 7);
match handler.run(&row(RETENTION_JOB_KIND)).await {
JobOutcome::Reschedule(delay) => assert_eq!(delay, DAILY),
other => panic!("{other:?}"),
}
assert!(Job::find_by_id("old", &database).await.unwrap().is_none());
}
#[tokio::test]
async fn recovery_queues_one_occurrence_however_often_it_runs() {
let (database, queue) = setup().await;
let handler = SweepJob::nonces(database.clone(), Duration::from_secs(300));
handler.recover(&queue).await;
handler.recover(&queue).await;
handler.recover(&queue).await;
assert_eq!(
Job::count_live(NONCE_SWEEP_KIND, &database).await.unwrap(),
1
);
}
#[tokio::test]
async fn no_target_ever_retires_itself_on_a_failure() {
let (database, _queue) = setup().await;
database.pool.close().await;
for handler in [
SweepJob::nonces(database.clone(), Duration::from_secs(300)),
SweepJob::audit(database.clone(), 7),
SweepJob::admin_sessions(database.clone(), Duration::from_secs(3600), 43200),
SweepJob::jobs(database.clone(), 7),
] {
let kind = handler.kind();
assert!(
matches!(handler.run(&row(kind)).await, JobOutcome::Reschedule(_)),
"{kind} must reschedule rather than retire"
);
}
}
}