use sqlx::PgPool;
use super::{terminal_state_labels, PurgeReport, RetentionPolicy};
pub async fn purge(
pool: &PgPool,
table: &'static str,
policy: &RetentionPolicy,
) -> Result<PurgeReport, sqlx::Error> {
let labels = terminal_state_labels();
let interval = format!("{} seconds", policy.terminal_max_age.as_secs());
let batch = i64::from(policy.effective_batch_size());
let sql = format!(
"DELETE FROM {table} WHERE ctid IN ( \
SELECT ctid FROM {table} \
WHERE state = ANY($1) \
AND updated_at < now() - $2::interval \
LIMIT $3 \
)"
);
let mut report = PurgeReport::default();
loop {
if policy.max_batches.is_some_and(|max| report.batches >= max) {
report.complete = false;
break;
}
let deleted = sqlx::query(&sql)
.bind(&labels)
.bind(&interval)
.bind(batch)
.execute(pool)
.await?
.rows_affected();
if deleted == 0 {
report.complete = true;
break;
}
report.tasks_deleted += deleted;
report.batches += 1;
}
Ok(report)
}