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";
pub const ORDER_SWEEP_KIND: &str = "order_sweep";
const DAILY: Duration = Duration::from_secs(24 * 60 * 60);
#[derive(Debug, Clone)]
enum SweepTarget {
Nonces { ttl: Duration },
AuditLog { retention_days: u64 },
AdminSessions { idle_timeout: Duration },
Jobs { retention_days: u64 },
Orders { retention: Vec<(String, 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 orders(database: Arc<Database>, retention: Vec<(String, u64)>) -> Self {
Self {
target: SweepTarget::Orders { retention },
interval: DAILY,
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);
}
}
SweepTarget::Orders { retention } => {
for (profile, retention_days) in retention {
let cutoff = crate::admin::ops::audit_cutoff(*retention_days);
match crate::sqlite::order::Order::cleanup(profile, cutoff, &self.database)
.await
{
Ok(removed) => info!(
event = "order_reaper_swept",
outcome = "success",
profile = %profile,
rows_removed = removed,
cutoff
),
Err(error) => {
error!(event = "order_reaper_failed", outcome = "failure", profile = %profile, 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,
SweepTarget::Orders { .. } => ORDER_SWEEP_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.clone(), 7).kind(),
SweepJob::orders(database, vec![("default".to_string(), 30)]).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),
SweepJob::orders(database.clone(), vec![("default".to_string(), 30)]),
] {
let kind = handler.kind();
assert!(
matches!(handler.run(&row(kind)).await, JobOutcome::Reschedule(_)),
"{kind} must reschedule rather than retire"
);
}
}
#[tokio::test]
async fn the_order_sweep_spares_valid_orders_and_takes_expired_ones() {
use crate::sqlite::account::Account;
use crate::sqlite::order::{Identifier, Order};
let (database, _queue) = setup().await;
let (account, _) = Account::find_or_create(
"default",
&[3u8; 8],
vec![],
&crate::audit::ClientContext::default(),
&database,
)
.await
.unwrap();
let ancient = now_secs() - 400 * 24 * 60 * 60;
let mut ids = Vec::new();
for _ in 0..3 {
let order = Order::create(
"default",
&account.id,
vec![Identifier::dns("example.com")],
ancient,
None,
None,
&database,
)
.await
.unwrap();
ids.push(order.id);
}
let (expired_pending, expired_valid, fresh) = (&ids[0], &ids[1], &ids[2]);
sqlx::query("UPDATE orders SET status = 'valid' WHERE id = ?;")
.bind(expired_valid)
.execute(&database.pool)
.await
.unwrap();
sqlx::query("UPDATE orders SET expires = ? WHERE id = ?;")
.bind(now_secs() + 3600)
.bind(fresh)
.execute(&database.pool)
.await
.unwrap();
let removed = Order::cleanup("default", crate::admin::ops::audit_cutoff(30), &database)
.await
.unwrap();
assert_eq!(removed, 1, "only the expired, undecided order goes");
assert!(
Order::find_by_id(expired_pending, &database)
.await
.unwrap()
.is_none()
);
assert!(
Order::find_by_id(expired_valid, &database)
.await
.unwrap()
.is_some(),
"a valid order is never swept: its row is how a certificate is revoked"
);
assert!(Order::find_by_id(fresh, &database).await.unwrap().is_some());
}
#[tokio::test]
async fn the_order_sweep_is_scoped_to_one_profile() {
use crate::sqlite::account::Account;
use crate::sqlite::order::{Identifier, Order};
let (database, _queue) = setup().await;
let ancient = now_secs() - 400 * 24 * 60 * 60;
let mut ids = Vec::new();
for (index, profile) in ["default", "other"].iter().enumerate() {
let (account, _) = Account::find_or_create(
profile,
&[index as u8 + 40; 8],
vec![],
&crate::audit::ClientContext::default(),
&database,
)
.await
.unwrap();
let order = Order::create(
profile,
&account.id,
vec![Identifier::dns("example.com")],
ancient,
None,
None,
&database,
)
.await
.unwrap();
ids.push(order.id);
}
let removed = Order::cleanup("default", crate::admin::ops::audit_cutoff(30), &database)
.await
.unwrap();
assert_eq!(removed, 1);
assert!(
Order::find_by_id(&ids[0], &database)
.await
.unwrap()
.is_none()
);
assert!(
Order::find_by_id(&ids[1], &database)
.await
.unwrap()
.is_some(),
"another profile's orders are not this profile's to delete"
);
}
}