use crate::local_outputs::{
LocalOutputFilter, LocalOutputRecord, SweepReport, SweepScope,
sweep::{self, SweepOptions},
};
use crate::serve::error::ServeError;
use crate::serve::rbac::{AuthContext, Permission};
use crate::serve::state::ServerState;
use axum::Json;
use axum::extract::{Extension, Path, Query, State};
use chrono::Utc;
use serde::{Deserialize, Serialize};
const DEFAULT_LIMIT: usize = 200;
const MAX_LIMIT: usize = 2000;
#[derive(Debug, Deserialize)]
pub struct ListQuery {
pub dataset_id: Option<String>,
pub pipeline: Option<String>,
#[serde(default)]
pub include_expired: bool,
pub limit: Option<usize>,
}
#[derive(Debug, Clone, Serialize)]
pub struct LocalOutputView {
#[serde(flatten)]
pub record: LocalOutputRecord,
pub state: &'static str,
pub age_secs: u64,
pub retention_days_effective: Option<u32>,
}
#[derive(Debug, Serialize)]
pub struct ListResponse {
pub outputs: Vec<LocalOutputView>,
pub retention_days: u32,
pub gc_enabled: bool,
pub can_manage: bool,
pub preview_enabled: bool,
pub preview_default_rows: Option<usize>,
pub preview_max_rows: Option<usize>,
}
pub async fn list_outputs(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Query(query): Query<ListQuery>,
) -> Result<Json<ListResponse>, ServeError> {
let retention_days = retention_days(&state);
let filter = LocalOutputFilter {
dataset_id: query.dataset_id,
pipeline: query.pipeline,
include_deleted: query.include_expired,
limit: query.limit.unwrap_or(DEFAULT_LIMIT).clamp(1, MAX_LIMIT),
};
let now = Utc::now();
let outputs = state
.history()
.local_output_list(&filter)
.await
.map_err(|e| ServeError::Internal(e.to_string()))?
.into_iter()
.map(|record| LocalOutputView {
state: record.state().as_str(),
age_secs: record.age_secs(now),
retention_days_effective: record.effective_retention_days(retention_days),
record,
})
.collect();
let preview = *state.preview();
Ok(Json(ListResponse {
outputs,
retention_days,
gc_enabled: retention_days > 0,
can_manage: actor.role.grants(Permission::LocalOutputManage),
preview_enabled: preview.enabled,
preview_default_rows: preview.default_rows(),
preview_max_rows: preview.max_rows(),
}))
}
pub async fn delete_output(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Path(id): Path<String>,
) -> Result<Json<SweepReport>, ServeError> {
if state
.history()
.local_output_get(&id)
.await
.map_err(|e| ServeError::Internal(e.to_string()))?
.is_none()
{
return Err(ServeError::NotFound);
}
let report = sweep::run(
state.history().as_ref(),
&SweepScope::Output(id),
&options(&state, false),
)
.await
.map_err(|e| ServeError::Internal(e.to_string()))?;
audit(&state, &actor, "local_output.delete", &report).await;
Ok(Json(report))
}
#[derive(Debug, Default, Deserialize)]
pub struct CleanupRequest {
#[serde(default)]
pub older_than_days: Option<u32>,
#[serde(default)]
pub expired: bool,
#[serde(default)]
pub all: bool,
#[serde(default)]
pub dataset_id: Option<String>,
#[serde(default)]
pub run_id: Option<String>,
#[serde(default)]
pub dry_run: bool,
#[serde(default)]
pub confirm: bool,
}
impl CleanupRequest {
fn scope(&self) -> Result<SweepScope, ServeError> {
let chosen: Vec<SweepScope> = [
self.older_than_days.map(SweepScope::OlderThanDays),
self.expired.then_some(SweepScope::Expired),
self.all.then_some(SweepScope::All),
self.dataset_id.clone().map(SweepScope::Dataset),
self.run_id.clone().map(SweepScope::Run),
]
.into_iter()
.flatten()
.collect();
match chosen.len() {
1 => Ok(chosen.into_iter().next().expect("length checked")),
0 => Err(ServeError::BadConfig(
"cleanup: choose a scope — one of `older_than_days`, `expired`, \
`dataset_id`, `run_id`, or `all`"
.into(),
)),
_ => Err(ServeError::BadConfig(
"cleanup: `older_than_days`, `expired`, `dataset_id`, `run_id`, and `all` \
are mutually exclusive — send exactly one"
.into(),
)),
}
}
}
pub async fn cleanup(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Json(req): Json<CleanupRequest>,
) -> Result<Json<SweepReport>, ServeError> {
let scope = req.scope()?;
if scope.requires_confirmation() && !req.confirm && !req.dry_run {
return Err(ServeError::BadConfig(format!(
"cleanup: scope `{}` deletes tracked outputs that are still inside their \
retention window — resend with `\"confirm\": true` to proceed, or \
`\"dry_run\": true` to see what it would remove",
scope.label()
)));
}
let report = sweep::run(
state.history().as_ref(),
&scope,
&options(&state, req.dry_run),
)
.await
.map_err(|e| ServeError::Internal(e.to_string()))?;
audit(&state, &actor, "local_output.cleanup", &report).await;
Ok(Json(report))
}
async fn audit(state: &ServerState, actor: &AuthContext, action: &str, report: &SweepReport) {
let result = if report.dry_run {
"dry_run".to_string()
} else {
format!("deleted={} skipped={}", report.deleted, report.skipped)
};
crate::serve::audit::write(state, actor, action, None, None, &result).await;
}
fn retention_days(state: &ServerState) -> u32 {
state.local_output_retention_days()
}
fn options(state: &ServerState, dry_run: bool) -> SweepOptions {
SweepOptions::new(retention_days(state))
.dry_run(dry_run)
.in_flight(state.registry().live_run_ids())
.in_flight_grace(state.local_output_in_flight_grace())
}
#[cfg(test)]
mod tests {
use super::*;
fn req() -> CleanupRequest {
CleanupRequest::default()
}
#[test]
fn each_field_resolves_to_its_scope() {
let mut r = req();
r.older_than_days = Some(3);
assert_eq!(r.scope().unwrap(), SweepScope::OlderThanDays(3));
let mut r = req();
r.expired = true;
assert_eq!(r.scope().unwrap(), SweepScope::Expired);
let mut r = req();
r.all = true;
assert_eq!(r.scope().unwrap(), SweepScope::All);
let mut r = req();
r.dataset_id = Some("ds1".into());
assert_eq!(r.scope().unwrap(), SweepScope::Dataset("ds1".into()));
let mut r = req();
r.run_id = Some("run-1".into());
assert_eq!(r.scope().unwrap(), SweepScope::Run("run-1".into()));
}
#[test]
fn an_empty_request_is_rejected_rather_than_defaulting() {
let err = req().scope().unwrap_err();
assert!(matches!(err, ServeError::BadConfig(_)), "{err:?}");
}
#[test]
fn combining_scopes_is_rejected() {
let mut r = req();
r.all = true;
r.older_than_days = Some(30);
let err = r.scope().unwrap_err();
match err {
ServeError::BadConfig(m) => assert!(m.contains("mutually exclusive"), "{m}"),
other => panic!("expected BadConfig, got {other:?}"),
}
}
#[test]
fn confirmation_is_required_for_exactly_the_unbounded_scopes() {
let mut r = req();
r.all = true;
assert!(r.scope().unwrap().requires_confirmation());
let mut r = req();
r.older_than_days = Some(0);
assert!(
r.scope().unwrap().requires_confirmation(),
"a zero-day window matches every row"
);
let mut r = req();
r.older_than_days = Some(7);
assert!(!r.scope().unwrap().requires_confirmation());
let mut r = req();
r.dataset_id = Some("ds".into());
assert!(!r.scope().unwrap().requires_confirmation());
}
#[test]
fn dry_run_is_independent_of_the_scope() {
let mut r = req();
r.all = true;
r.dry_run = true;
assert_eq!(r.scope().unwrap(), SweepScope::All);
assert!(r.dry_run);
}
#[test]
fn request_parses_from_the_console_payloads() {
let r: CleanupRequest = serde_json::from_str(r#"{"older_than_days": 7}"#).unwrap();
assert_eq!(r.scope().unwrap(), SweepScope::OlderThanDays(7));
let r: CleanupRequest = serde_json::from_str(r#"{"all": true}"#).unwrap();
assert_eq!(r.scope().unwrap(), SweepScope::All);
let r: CleanupRequest =
serde_json::from_str(r#"{"dataset_id": "abc", "dry_run": true}"#).unwrap();
assert_eq!(r.scope().unwrap(), SweepScope::Dataset("abc".into()));
assert!(r.dry_run);
}
}