feat(jobs): add admin call to purge jubs result
This commit is contained in:
@@ -431,6 +431,30 @@ impl JobStoreProvider for PgJobStoreProvider {
|
||||
.collect())
|
||||
}
|
||||
|
||||
async fn purge_terminal_runs(&self, retention_days: i32) -> Result<u64, DomainError> {
|
||||
// Defensive floor — zero would eat just-completed runs;
|
||||
// negative would eat the whole terminal history.
|
||||
let days = retention_days.max(1);
|
||||
// `ON DELETE CASCADE` on jobs.run_findings.run_id
|
||||
// (migration 20260930000001) drops findings with their
|
||||
// parent run. Non-terminal statuses (Running / Paused /
|
||||
// CancelRequested) explicitly excluded to protect
|
||||
// in-flight work.
|
||||
let result = sqlx::query(
|
||||
r#"
|
||||
DELETE FROM jobs.recoverable_runs
|
||||
WHERE status IN ('Completed', 'Failed')
|
||||
AND completed_at IS NOT NULL
|
||||
AND completed_at < NOW() - ($1 || ' days')::interval
|
||||
"#,
|
||||
)
|
||||
.bind(days.to_string())
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| map_sqlx_err("purge_terminal_runs", e))?;
|
||||
Ok(result.rows_affected())
|
||||
}
|
||||
|
||||
async fn request_cancel(&self, job_name: &str) -> Result<Option<Uuid>, DomainError> {
|
||||
// Only Running → CancelRequested flips. `Paused` can be
|
||||
// cancelled by not resuming — no need for a state change.
|
||||
|
||||
@@ -371,6 +371,29 @@ pub trait JobStoreProvider: Send + Sync {
|
||||
&self,
|
||||
run_id: Uuid,
|
||||
) -> Result<Vec<(String, u64)>, DomainError>;
|
||||
|
||||
/// Operator-triggered retention cleanup. DELETEs every
|
||||
/// TERMINAL run (`Completed`, `Failed`) whose `completed_at`
|
||||
/// is older than `retention_days` days ago. Findings drop
|
||||
/// alongside via the `ON DELETE CASCADE` FK on
|
||||
/// `jobs.run_findings.run_id`.
|
||||
///
|
||||
/// Non-terminal rows (`Running`, `Paused`, `CancelRequested`)
|
||||
/// are ALWAYS preserved regardless of age — an in-flight or
|
||||
/// paused run must not be reaped by retention.
|
||||
///
|
||||
/// `retention_days` is treated as `max(1, retention_days)` at
|
||||
/// the impl layer to defend against a zero/negative value
|
||||
/// eating just-completed runs.
|
||||
///
|
||||
/// Returns the number of run rows deleted (which equals
|
||||
/// the number of finding rows deleted *transitively* via
|
||||
/// CASCADE; callers wanting the finding count separately
|
||||
/// should query it BEFORE calling this).
|
||||
///
|
||||
/// Powers `POST /api/admin/jobs/runs/purge`. Not periodic — the
|
||||
/// operator decides when to reclaim space.
|
||||
async fn purge_terminal_runs(&self, retention_days: i32) -> Result<u64, DomainError>;
|
||||
}
|
||||
|
||||
/// How a `RunProgress` fraction was derived. Lets the UI communicate
|
||||
@@ -1094,6 +1117,25 @@ mod tests {
|
||||
Ok(counts.into_iter().collect())
|
||||
}
|
||||
|
||||
async fn purge_terminal_runs(&self, retention_days: i32) -> Result<u64, DomainError> {
|
||||
// Test-double: no `completed_at` to compare against, so
|
||||
// just drop every terminal-state store when
|
||||
// `retention_days` > 0. Sufficient for the trait
|
||||
// contract check; PG impl exercises the real
|
||||
// `completed_at < NOW() - days` filter.
|
||||
let days = retention_days.max(1);
|
||||
if days == 0 {
|
||||
return Ok(0);
|
||||
}
|
||||
let mut stores = self.stores.lock().unwrap();
|
||||
let before = stores.len();
|
||||
stores.retain(|s| {
|
||||
let state = s.state.lock().unwrap();
|
||||
!matches!(state.status, RunStatus::Completed | RunStatus::Failed)
|
||||
});
|
||||
Ok((before - stores.len()) as u64)
|
||||
}
|
||||
|
||||
async fn request_cancel(&self, _job_name: &str) -> Result<Option<Uuid>, DomainError> {
|
||||
let stores = self.stores.lock().unwrap();
|
||||
if let Some(s) = stores.last() {
|
||||
|
||||
@@ -164,6 +164,9 @@ pub fn admin_routes() -> Router<Arc<AppState>> {
|
||||
"/jobs/{name}/runs/{id}/findings",
|
||||
get(list_job_run_findings),
|
||||
)
|
||||
// Retention cleanup — operator-triggered, not periodic.
|
||||
// See `purge_job_runs` docstring for the semantics.
|
||||
.route("/jobs/runs/purge", post(purge_job_runs))
|
||||
// Drives — admin-wide view (distinct from `/api/drives` which
|
||||
// is filtered to the caller's role grants).
|
||||
.route("/drives", get(list_all_drives))
|
||||
@@ -2392,3 +2395,74 @@ pub async fn list_job_run_findings(
|
||||
Err(e) => AppError::internal_error(format!("list_findings failed: {e}")).into_response(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Query parameters for `POST /api/admin/jobs/runs/purge`.
|
||||
///
|
||||
/// `days` — retention window. Terminal runs (`Completed`, `Failed`)
|
||||
/// with `completed_at` older than this many days ago are deleted
|
||||
/// (with their findings via CASCADE). Default 30. Minimum enforced
|
||||
/// at 1 by the provider — zero would eat runs completed seconds
|
||||
/// ago. Non-terminal runs are ALWAYS preserved regardless of age.
|
||||
#[derive(serde::Deserialize)]
|
||||
pub struct PurgeJobRunsQuery {
|
||||
#[serde(default = "default_purge_days")]
|
||||
pub days: i32,
|
||||
}
|
||||
|
||||
fn default_purge_days() -> i32 {
|
||||
30
|
||||
}
|
||||
|
||||
/// `POST /api/admin/jobs/runs/purge?days=N` — operator-triggered
|
||||
/// cleanup of old terminal runs + their findings. Not periodic;
|
||||
/// admins fire this when they want to reclaim `jobs.*` history
|
||||
/// space. Delegates entirely to
|
||||
/// `JobStoreProvider::purge_terminal_runs` — no SQL in the handler
|
||||
/// (see `AGENTS.md` § handler thinness).
|
||||
#[utoipa::path(
|
||||
post,
|
||||
path = "/api/admin/jobs/runs/purge",
|
||||
params(
|
||||
("days" = Option<i32>, Query, description = "Retention window in days (default 30, minimum 1). Terminal runs older than this are deleted with their findings; non-terminal runs are always preserved."),
|
||||
),
|
||||
responses(
|
||||
(status = 200, description = "Purge complete"),
|
||||
(status = 401, description = "Unauthorized"),
|
||||
(status = 403, description = "Admin required"),
|
||||
(status = 500, description = "DB error"),
|
||||
),
|
||||
security(("bearerAuth" = [])),
|
||||
tag = "admin"
|
||||
)]
|
||||
pub async fn purge_job_runs(
|
||||
State(state): State<Arc<AppState>>,
|
||||
axum::extract::Query(query): axum::extract::Query<PurgeJobRunsQuery>,
|
||||
) -> impl IntoResponse {
|
||||
use crate::infrastructure::scheduler::JobStoreProvider as _;
|
||||
let retention_days = query.days.max(1);
|
||||
match state
|
||||
.core
|
||||
.job_store_provider
|
||||
.purge_terminal_runs(retention_days)
|
||||
.await
|
||||
{
|
||||
Ok(purged) => {
|
||||
tracing::info!(
|
||||
target: "audit",
|
||||
event = "jobs.runs_purged",
|
||||
purged = purged,
|
||||
retention_days = retention_days,
|
||||
"👮🏻♂️ admin purged {purged} terminal recoverable-run row(s) past {retention_days} day retention (findings cascaded)",
|
||||
);
|
||||
(
|
||||
StatusCode::OK,
|
||||
Json(serde_json::json!({
|
||||
"purged": purged,
|
||||
"retention_days": retention_days,
|
||||
})),
|
||||
)
|
||||
.into_response()
|
||||
}
|
||||
Err(e) => AppError::internal_error(format!("purge failed: {e}")).into_response(),
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user