mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-06 02:48:25 +00:00
Improve arc system prune (#14)
* arc(01KKA1A9ZHFMWM3T9VC6HBWJ19): toolchain (success) Arc-Run: 01KKA1A9ZHFMWM3T9VC6HBWJ19 Arc-Completed: 2 Arc-Checkpoint: a38cd3755e5a5ad6a406e3557665c3ea6ba06cc7 * arc(01KKA1A9ZHFMWM3T9VC6HBWJ19): preflight_compile (success) Arc-Run: 01KKA1A9ZHFMWM3T9VC6HBWJ19 Arc-Completed: 3 Arc-Checkpoint: 62150d0142e04fad46092b85e5942ac1979da51b * arc(01KKA1A9ZHFMWM3T9VC6HBWJ19): preflight_lint (fail) Arc-Run: 01KKA1A9ZHFMWM3T9VC6HBWJ19 Arc-Completed: 4 Arc-Checkpoint: 2a213ceb22392e2dd1e00c4e66912ed57f0972af * arc(01KKA1A9ZHFMWM3T9VC6HBWJ19): fix_lints (success) Arc-Run: 01KKA1A9ZHFMWM3T9VC6HBWJ19 Arc-Completed: 5 Arc-Checkpoint: e9d720a9e2288664ebea475c4a1445532fde3ed9 * arc(01KKA1A9ZHFMWM3T9VC6HBWJ19): preflight_lint (success) Arc-Run: 01KKA1A9ZHFMWM3T9VC6HBWJ19 Arc-Completed: 6 Arc-Checkpoint: 8ef9490d981389c374970e0374ef48198fb19829 * arc(01KKA1A9ZHFMWM3T9VC6HBWJ19): implement (success) Arc-Run: 01KKA1A9ZHFMWM3T9VC6HBWJ19 Arc-Completed: 7 Arc-Checkpoint: 063ba4fe9554c6e0d16aaad40c4ffc7badbd0511 * arc(01KKA1A9ZHFMWM3T9VC6HBWJ19): simplify (success) Arc-Run: 01KKA1A9ZHFMWM3T9VC6HBWJ19 Arc-Completed: 8 Arc-Checkpoint: 4c1319c2a2d7c818faa460413e7993c93d30750f * arc(01KKA1A9ZHFMWM3T9VC6HBWJ19): verify (success) Arc-Run: 01KKA1A9ZHFMWM3T9VC6HBWJ19 Arc-Completed: 9 Arc-Checkpoint: a01dd1e93b231415be85a1925f90b6fedf1844fc --------- Co-authored-by: arc <arc@local>
This commit is contained in:
parent
5fb3fd485e
commit
d432d370c0
1 changed files with 498 additions and 34 deletions
|
|
@ -1,11 +1,40 @@
|
|||
use std::collections::HashMap;
|
||||
use std::fmt;
|
||||
use std::path::{Path, PathBuf};
|
||||
|
||||
use anyhow::{bail, Context, Result};
|
||||
use chrono::{DateTime, Utc};
|
||||
use clap::Args;
|
||||
use serde::Serialize;
|
||||
use tracing::{debug, info, warn};
|
||||
|
||||
use crate::outcome::StageStatus;
|
||||
|
||||
/// Status of a run directory — either concluded with a `StageStatus`, actively running, or unknown.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum RunStatus {
|
||||
Concluded(StageStatus),
|
||||
Running,
|
||||
Unknown,
|
||||
}
|
||||
|
||||
impl RunStatus {
|
||||
pub fn is_running(&self) -> bool {
|
||||
matches!(self, RunStatus::Running)
|
||||
}
|
||||
}
|
||||
|
||||
impl fmt::Display for RunStatus {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
match self {
|
||||
RunStatus::Concluded(s) => write!(f, "{s}"),
|
||||
RunStatus::Running => write!(f, "running"),
|
||||
RunStatus::Unknown => write!(f, "unknown"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Args)]
|
||||
pub struct RunFilterArgs {
|
||||
/// Only include runs started before this date (YYYY-MM-DD prefix match)
|
||||
|
|
@ -40,6 +69,10 @@ pub struct RunsPruneArgs {
|
|||
#[command(flatten)]
|
||||
pub filter: RunFilterArgs,
|
||||
|
||||
/// Only prune runs older than this duration (e.g. 24h, 7d). Default: 24h when no explicit filters are set.
|
||||
#[arg(long, value_name = "DURATION", value_parser = parse_duration)]
|
||||
pub older_than: Option<chrono::Duration>,
|
||||
|
||||
/// Actually delete (default is dry-run)
|
||||
#[arg(long)]
|
||||
pub yes: bool,
|
||||
|
|
@ -50,10 +83,14 @@ pub struct RunInfo {
|
|||
pub run_id: String,
|
||||
pub dir_name: String,
|
||||
pub workflow_name: String,
|
||||
pub status: String,
|
||||
pub status: RunStatus,
|
||||
pub start_time: String,
|
||||
pub labels: HashMap<String, String>,
|
||||
#[serde(skip)]
|
||||
pub start_time_dt: Option<DateTime<Utc>>,
|
||||
#[serde(skip)]
|
||||
pub end_time: Option<DateTime<Utc>>,
|
||||
#[serde(skip)]
|
||||
pub path: PathBuf,
|
||||
#[serde(skip)]
|
||||
pub is_orphan: bool,
|
||||
|
|
@ -85,10 +122,11 @@ pub fn scan_runs(base: &Path) -> Result<Vec<RunInfo>> {
|
|||
|
||||
let run_id = manifest.run_id;
|
||||
let workflow_name = manifest.workflow_name;
|
||||
let start_time = manifest.start_time.to_rfc3339();
|
||||
let start_time_dt = manifest.start_time;
|
||||
let start_time = start_time_dt.to_rfc3339();
|
||||
let labels = manifest.labels;
|
||||
|
||||
let status = read_status(&path);
|
||||
let (status, end_time) = read_status(&path);
|
||||
|
||||
runs.push(RunInfo {
|
||||
run_id,
|
||||
|
|
@ -97,28 +135,31 @@ pub fn scan_runs(base: &Path) -> Result<Vec<RunInfo>> {
|
|||
status,
|
||||
start_time,
|
||||
labels,
|
||||
start_time_dt: Some(start_time_dt),
|
||||
end_time,
|
||||
path,
|
||||
is_orphan: false,
|
||||
});
|
||||
} else {
|
||||
// Orphan directory — no manifest.json
|
||||
let mtime = entry
|
||||
let mtime_dt = entry
|
||||
.metadata()
|
||||
.ok()
|
||||
.and_then(|m| m.modified().ok())
|
||||
.map(|t| {
|
||||
let dt: chrono::DateTime<chrono::Utc> = t.into();
|
||||
dt.to_rfc3339()
|
||||
})
|
||||
.map(|t| -> DateTime<Utc> { t.into() });
|
||||
let mtime = mtime_dt
|
||||
.map(|dt| dt.to_rfc3339())
|
||||
.unwrap_or_default();
|
||||
|
||||
runs.push(RunInfo {
|
||||
run_id: dir_name.clone(),
|
||||
dir_name,
|
||||
workflow_name: "[no manifest]".to_string(),
|
||||
status: "unknown".to_string(),
|
||||
status: RunStatus::Unknown,
|
||||
start_time: mtime,
|
||||
labels: HashMap::new(),
|
||||
start_time_dt: mtime_dt,
|
||||
end_time: None,
|
||||
path,
|
||||
is_orphan: true,
|
||||
});
|
||||
|
|
@ -130,14 +171,14 @@ pub fn scan_runs(base: &Path) -> Result<Vec<RunInfo>> {
|
|||
Ok(runs)
|
||||
}
|
||||
|
||||
fn read_status(run_dir: &Path) -> String {
|
||||
fn read_status(run_dir: &Path) -> (RunStatus, Option<DateTime<Utc>>) {
|
||||
if let Ok(conclusion) = crate::conclusion::Conclusion::load(&run_dir.join("conclusion.json")) {
|
||||
return conclusion.status.to_string();
|
||||
return (RunStatus::Concluded(conclusion.status), Some(conclusion.timestamp));
|
||||
}
|
||||
if run_dir.join("run.pid").exists() {
|
||||
return "running".to_string();
|
||||
return (RunStatus::Running, None);
|
||||
}
|
||||
"unknown".to_string()
|
||||
(RunStatus::Unknown, None)
|
||||
}
|
||||
|
||||
/// Filter runs by criteria. Orphans are excluded unless `include_orphans` is true.
|
||||
|
|
@ -326,8 +367,8 @@ pub fn df_from(args: &DfArgs, data_dir: &Path, logs_base: &Path) -> Result<()> {
|
|||
struct RunSizeInfo {
|
||||
run_id: String,
|
||||
workflow_name: String,
|
||||
status: String,
|
||||
start_time: String,
|
||||
status: RunStatus,
|
||||
start_time_dt: Option<DateTime<Utc>>,
|
||||
size: u64,
|
||||
}
|
||||
|
||||
|
|
@ -336,7 +377,7 @@ pub fn df_from(args: &DfArgs, data_dir: &Path, logs_base: &Path) -> Result<()> {
|
|||
for run in &runs {
|
||||
let size = dir_size(&run.path);
|
||||
total_run_size += size;
|
||||
let is_active = run.status == "running";
|
||||
let is_active = run.status.is_running();
|
||||
if is_active {
|
||||
active_count += 1;
|
||||
} else {
|
||||
|
|
@ -347,7 +388,7 @@ pub fn df_from(args: &DfArgs, data_dir: &Path, logs_base: &Path) -> Result<()> {
|
|||
run_id: run.run_id.clone(),
|
||||
workflow_name: run.workflow_name.clone(),
|
||||
status: run.status.clone(),
|
||||
start_time: run.start_time.clone(),
|
||||
start_time_dt: run.start_time_dt,
|
||||
size,
|
||||
});
|
||||
}
|
||||
|
|
@ -450,7 +491,7 @@ pub fn df_from(args: &DfArgs, data_dir: &Path, logs_base: &Path) -> Result<()> {
|
|||
} else {
|
||||
detail.workflow_name.clone()
|
||||
};
|
||||
let age = if let Ok(dt) = chrono::DateTime::parse_from_rfc3339(&detail.start_time) {
|
||||
let age = if let Some(dt) = detail.start_time_dt {
|
||||
let dur = now.signed_duration_since(dt);
|
||||
if dur.num_days() > 0 {
|
||||
format!("{}d", dur.num_days())
|
||||
|
|
@ -462,7 +503,7 @@ pub fn df_from(args: &DfArgs, data_dir: &Path, logs_base: &Path) -> Result<()> {
|
|||
} else {
|
||||
"-".to_string()
|
||||
};
|
||||
let reclaimable_marker = if detail.status != "running" { " *" } else { "" };
|
||||
let reclaimable_marker = if !detail.status.is_running() { " *" } else { "" };
|
||||
println!(
|
||||
"{:<30} {:<18} {:<10} {:>5} {:>10}{}",
|
||||
run_id_display,
|
||||
|
|
@ -485,10 +526,27 @@ pub fn prune_command(args: &RunsPruneArgs) -> Result<()> {
|
|||
prune_from(args, &base)
|
||||
}
|
||||
|
||||
/// Parse a human duration string like "24h" or "7d" into a `chrono::Duration`.
|
||||
fn parse_duration(s: &str) -> Result<chrono::Duration> {
|
||||
let s = s.trim();
|
||||
if s.is_empty() {
|
||||
bail!("empty duration string");
|
||||
}
|
||||
let (num_str, unit) = s.split_at(s.len() - 1);
|
||||
let num: i64 = num_str
|
||||
.parse()
|
||||
.with_context(|| format!("invalid duration: {s}"))?;
|
||||
match unit {
|
||||
"h" => Ok(chrono::Duration::hours(num)),
|
||||
"d" => Ok(chrono::Duration::days(num)),
|
||||
_ => bail!("invalid duration unit '{unit}' in '{s}' (expected 'h' or 'd')"),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn prune_from(args: &RunsPruneArgs, base: &Path) -> Result<()> {
|
||||
let runs = scan_runs(base)?;
|
||||
let label_filters = parse_label_filters(&args.filter.label);
|
||||
let filtered = filter_runs(
|
||||
let mut filtered = filter_runs(
|
||||
&runs,
|
||||
args.filter.before.as_deref(),
|
||||
args.filter.workflow.as_deref(),
|
||||
|
|
@ -496,25 +554,66 @@ pub fn prune_from(args: &RunsPruneArgs, base: &Path) -> Result<()> {
|
|||
args.filter.orphans,
|
||||
);
|
||||
|
||||
// Determine if the user passed any explicit filters
|
||||
let has_explicit_filters = args.filter.before.is_some()
|
||||
|| args.filter.workflow.is_some()
|
||||
|| !args.filter.label.is_empty()
|
||||
|| args.filter.orphans;
|
||||
|
||||
// Apply staleness filter: default 24h when no explicit filters, or use --older-than
|
||||
let staleness_threshold = if let Some(dur) = args.older_than {
|
||||
Some(dur)
|
||||
} else if !has_explicit_filters {
|
||||
Some(chrono::Duration::hours(24))
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
if let Some(threshold) = staleness_threshold {
|
||||
let now = Utc::now();
|
||||
let cutoff = now - threshold;
|
||||
filtered.retain(|run| {
|
||||
// Exclude running runs
|
||||
if run.status.is_running() {
|
||||
return false;
|
||||
}
|
||||
// Use end_time if available, fall back to start_time
|
||||
let effective_time = run.end_time.or(run.start_time_dt);
|
||||
match effective_time {
|
||||
Some(t) => t < cutoff,
|
||||
None => false,
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
if filtered.is_empty() {
|
||||
eprintln!("No matching runs to prune.");
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
// Calculate total disk space
|
||||
let total_bytes: u64 = filtered.iter().map(|r| dir_size(&r.path)).sum();
|
||||
info!(count = filtered.len(), bytes = total_bytes, "pruning runs");
|
||||
|
||||
if args.yes {
|
||||
for run in &filtered {
|
||||
info!(run_id = %run.run_id, path = %run.path.display(), "deleting run");
|
||||
std::fs::remove_dir_all(&run.path)?;
|
||||
}
|
||||
eprintln!("{} run(s) deleted.", filtered.len());
|
||||
eprintln!(
|
||||
"{} run(s) deleted ({} freed).",
|
||||
filtered.len(),
|
||||
format_size(total_bytes)
|
||||
);
|
||||
} else {
|
||||
for run in &filtered {
|
||||
debug!(run_id = %run.run_id, "would delete run (dry-run)");
|
||||
println!("would delete: {} ({})", run.dir_name, run.workflow_name);
|
||||
}
|
||||
eprintln!(
|
||||
"\n{} run(s) would be deleted. Pass --yes to confirm.",
|
||||
filtered.len()
|
||||
"\n{} run(s) would be deleted ({} freed). Pass --yes to confirm.",
|
||||
filtered.len(),
|
||||
format_size(total_bytes)
|
||||
);
|
||||
}
|
||||
Ok(())
|
||||
|
|
@ -584,13 +683,13 @@ mod tests {
|
|||
|
||||
let completed = runs.iter().find(|r| r.run_id == "abc123").unwrap();
|
||||
assert_eq!(completed.workflow_name, "my-pipeline");
|
||||
assert_eq!(completed.status, "success");
|
||||
assert_eq!(completed.status, RunStatus::Concluded(crate::outcome::StageStatus::Success));
|
||||
assert_eq!(completed.labels.get("env").unwrap(), "prod");
|
||||
assert!(!completed.is_orphan);
|
||||
|
||||
let orphan = runs.iter().find(|r| r.is_orphan).unwrap();
|
||||
assert_eq!(orphan.workflow_name, "[no manifest]");
|
||||
assert_eq!(orphan.status, "unknown");
|
||||
assert_eq!(orphan.status, RunStatus::Unknown);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
@ -615,7 +714,7 @@ mod tests {
|
|||
|
||||
let runs = scan_runs(base).unwrap();
|
||||
assert_eq!(runs.len(), 1);
|
||||
assert_eq!(runs[0].status, "running");
|
||||
assert_eq!(runs[0].status, RunStatus::Running);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
@ -638,9 +737,11 @@ mod tests {
|
|||
run_id: "old".into(),
|
||||
dir_name: "d1".into(),
|
||||
workflow_name: "p".into(),
|
||||
status: "success".into(),
|
||||
status: RunStatus::Concluded(crate::outcome::StageStatus::Success),
|
||||
start_time: "2025-06-01T00:00:00Z".into(),
|
||||
labels: HashMap::new(),
|
||||
start_time_dt: None,
|
||||
end_time: None,
|
||||
path: PathBuf::from("/tmp/d1"),
|
||||
is_orphan: false,
|
||||
},
|
||||
|
|
@ -648,9 +749,11 @@ mod tests {
|
|||
run_id: "new".into(),
|
||||
dir_name: "d2".into(),
|
||||
workflow_name: "p".into(),
|
||||
status: "success".into(),
|
||||
status: RunStatus::Concluded(crate::outcome::StageStatus::Success),
|
||||
start_time: "2026-03-01T00:00:00Z".into(),
|
||||
labels: HashMap::new(),
|
||||
start_time_dt: None,
|
||||
end_time: None,
|
||||
path: PathBuf::from("/tmp/d2"),
|
||||
is_orphan: false,
|
||||
},
|
||||
|
|
@ -667,9 +770,11 @@ mod tests {
|
|||
run_id: "a".into(),
|
||||
dir_name: "d1".into(),
|
||||
workflow_name: "deploy-prod".into(),
|
||||
status: "success".into(),
|
||||
status: RunStatus::Concluded(crate::outcome::StageStatus::Success),
|
||||
start_time: "2026-01-01T00:00:00Z".into(),
|
||||
labels: HashMap::new(),
|
||||
start_time_dt: None,
|
||||
end_time: None,
|
||||
path: PathBuf::from("/tmp/d1"),
|
||||
is_orphan: false,
|
||||
},
|
||||
|
|
@ -677,9 +782,11 @@ mod tests {
|
|||
run_id: "b".into(),
|
||||
dir_name: "d2".into(),
|
||||
workflow_name: "test-suite".into(),
|
||||
status: "success".into(),
|
||||
status: RunStatus::Concluded(crate::outcome::StageStatus::Success),
|
||||
start_time: "2026-01-01T00:00:00Z".into(),
|
||||
labels: HashMap::new(),
|
||||
start_time_dt: None,
|
||||
end_time: None,
|
||||
path: PathBuf::from("/tmp/d2"),
|
||||
is_orphan: false,
|
||||
},
|
||||
|
|
@ -696,9 +803,11 @@ mod tests {
|
|||
run_id: "a".into(),
|
||||
dir_name: "d1".into(),
|
||||
workflow_name: "p".into(),
|
||||
status: "success".into(),
|
||||
status: RunStatus::Concluded(crate::outcome::StageStatus::Success),
|
||||
start_time: "2026-01-01T00:00:00Z".into(),
|
||||
labels: HashMap::from([("env".into(), "prod".into())]),
|
||||
start_time_dt: None,
|
||||
end_time: None,
|
||||
path: PathBuf::from("/tmp/d1"),
|
||||
is_orphan: false,
|
||||
},
|
||||
|
|
@ -706,9 +815,11 @@ mod tests {
|
|||
run_id: "b".into(),
|
||||
dir_name: "d2".into(),
|
||||
workflow_name: "p".into(),
|
||||
status: "success".into(),
|
||||
status: RunStatus::Concluded(crate::outcome::StageStatus::Success),
|
||||
start_time: "2026-01-01T00:00:00Z".into(),
|
||||
labels: HashMap::from([("env".into(), "staging".into())]),
|
||||
start_time_dt: None,
|
||||
end_time: None,
|
||||
path: PathBuf::from("/tmp/d2"),
|
||||
is_orphan: false,
|
||||
},
|
||||
|
|
@ -730,9 +841,11 @@ mod tests {
|
|||
run_id: "orphan".into(),
|
||||
dir_name: "d1".into(),
|
||||
workflow_name: "[no manifest]".into(),
|
||||
status: "unknown".into(),
|
||||
status: RunStatus::Unknown,
|
||||
start_time: "".into(),
|
||||
labels: HashMap::new(),
|
||||
start_time_dt: None,
|
||||
end_time: None,
|
||||
path: PathBuf::from("/tmp/d1"),
|
||||
is_orphan: true,
|
||||
}];
|
||||
|
|
@ -772,6 +885,7 @@ mod tests {
|
|||
label: Vec::new(),
|
||||
orphans: false,
|
||||
},
|
||||
older_than: None,
|
||||
yes: false,
|
||||
};
|
||||
|
||||
|
|
@ -826,6 +940,7 @@ mod tests {
|
|||
label: Vec::new(),
|
||||
orphans: false,
|
||||
},
|
||||
older_than: None,
|
||||
yes: true,
|
||||
};
|
||||
|
||||
|
|
@ -851,6 +966,7 @@ mod tests {
|
|||
label: Vec::new(),
|
||||
orphans: true,
|
||||
},
|
||||
older_than: None,
|
||||
yes: true,
|
||||
};
|
||||
|
||||
|
|
@ -1056,4 +1172,352 @@ mod tests {
|
|||
"Should mention ambiguity"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// === Step 1: end_time tests ===
|
||||
|
||||
#[test]
|
||||
fn scan_runs_populates_end_time() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let base = tmp.path();
|
||||
|
||||
make_run_dir(
|
||||
base,
|
||||
"20260101-DONE",
|
||||
Some(serde_json::json!({
|
||||
"run_id": "done-1",
|
||||
"workflow_name": "pipeline",
|
||||
"goal": "",
|
||||
"start_time": "2026-01-01T12:00:00Z",
|
||||
"node_count": 1,
|
||||
"edge_count": 0
|
||||
})),
|
||||
Some(serde_json::json!({
|
||||
"timestamp": "2026-01-01T12:05:00Z",
|
||||
"status": "success",
|
||||
"duration_ms": 300000
|
||||
})),
|
||||
false,
|
||||
);
|
||||
|
||||
let runs = scan_runs(base).unwrap();
|
||||
assert_eq!(runs.len(), 1);
|
||||
let run = &runs[0];
|
||||
assert!(run.end_time.is_some());
|
||||
assert_eq!(
|
||||
run.end_time.unwrap().to_rfc3339(),
|
||||
"2026-01-01T12:05:00+00:00"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scan_runs_end_time_none_when_running() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let base = tmp.path();
|
||||
|
||||
make_run_dir(
|
||||
base,
|
||||
"20260101-RUNNING",
|
||||
Some(serde_json::json!({
|
||||
"run_id": "running-1",
|
||||
"workflow_name": "pipeline",
|
||||
"goal": "",
|
||||
"start_time": "2026-01-01T12:00:00Z",
|
||||
"node_count": 1,
|
||||
"edge_count": 0
|
||||
})),
|
||||
None,
|
||||
true,
|
||||
);
|
||||
|
||||
let runs = scan_runs(base).unwrap();
|
||||
assert_eq!(runs.len(), 1);
|
||||
assert_eq!(runs[0].status, RunStatus::Running);
|
||||
assert!(runs[0].end_time.is_none());
|
||||
}
|
||||
|
||||
// === Step 2: staleness heuristic tests ===
|
||||
|
||||
#[test]
|
||||
fn prune_default_skips_recent_runs() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let base = tmp.path();
|
||||
|
||||
// A run completed just now — should NOT be pruned by default
|
||||
let now = Utc::now();
|
||||
let recent_ts = now.to_rfc3339();
|
||||
let dir = make_run_dir(
|
||||
base,
|
||||
"20260309-RECENT",
|
||||
Some(serde_json::json!({
|
||||
"run_id": "recent-1",
|
||||
"workflow_name": "pipeline",
|
||||
"goal": "",
|
||||
"start_time": &recent_ts,
|
||||
"node_count": 1,
|
||||
"edge_count": 0
|
||||
})),
|
||||
Some(serde_json::json!({
|
||||
"timestamp": &recent_ts,
|
||||
"status": "success",
|
||||
"duration_ms": 1000
|
||||
})),
|
||||
false,
|
||||
);
|
||||
|
||||
let args = RunsPruneArgs {
|
||||
filter: RunFilterArgs {
|
||||
before: None,
|
||||
workflow: None,
|
||||
label: Vec::new(),
|
||||
orphans: false,
|
||||
},
|
||||
older_than: None,
|
||||
yes: false,
|
||||
};
|
||||
|
||||
prune_from(&args, base).unwrap();
|
||||
assert!(dir.exists(), "recent run should not be pruned by default");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn prune_default_skips_running_runs() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let base = tmp.path();
|
||||
|
||||
let old_ts = (Utc::now() - chrono::Duration::hours(48)).to_rfc3339();
|
||||
let dir = make_run_dir(
|
||||
base,
|
||||
"20260307-RUNNING",
|
||||
Some(serde_json::json!({
|
||||
"run_id": "running-old",
|
||||
"workflow_name": "pipeline",
|
||||
"goal": "",
|
||||
"start_time": &old_ts,
|
||||
"node_count": 1,
|
||||
"edge_count": 0
|
||||
})),
|
||||
None,
|
||||
true, // running
|
||||
);
|
||||
|
||||
let args = RunsPruneArgs {
|
||||
filter: RunFilterArgs {
|
||||
before: None,
|
||||
workflow: None,
|
||||
label: Vec::new(),
|
||||
orphans: false,
|
||||
},
|
||||
older_than: None,
|
||||
yes: false,
|
||||
};
|
||||
|
||||
prune_from(&args, base).unwrap();
|
||||
assert!(dir.exists(), "running run should not be pruned");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn prune_default_targets_stale_runs() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let base = tmp.path();
|
||||
|
||||
let old_ts = (Utc::now() - chrono::Duration::hours(48)).to_rfc3339();
|
||||
let dir = make_run_dir(
|
||||
base,
|
||||
"20260307-STALE",
|
||||
Some(serde_json::json!({
|
||||
"run_id": "stale-1",
|
||||
"workflow_name": "old-pipeline",
|
||||
"goal": "",
|
||||
"start_time": &old_ts,
|
||||
"node_count": 1,
|
||||
"edge_count": 0
|
||||
})),
|
||||
Some(serde_json::json!({
|
||||
"timestamp": &old_ts,
|
||||
"status": "success",
|
||||
"duration_ms": 1000
|
||||
})),
|
||||
false,
|
||||
);
|
||||
|
||||
let args = RunsPruneArgs {
|
||||
filter: RunFilterArgs {
|
||||
before: None,
|
||||
workflow: None,
|
||||
label: Vec::new(),
|
||||
orphans: false,
|
||||
},
|
||||
older_than: None,
|
||||
yes: true,
|
||||
};
|
||||
|
||||
prune_from(&args, base).unwrap();
|
||||
assert!(!dir.exists(), "stale run (48h old) should be pruned by default");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn prune_explicit_filter_bypasses_staleness() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let base = tmp.path();
|
||||
|
||||
// A run completed just now, but matched by --before
|
||||
let now = Utc::now();
|
||||
let recent_ts = now.to_rfc3339();
|
||||
let dir = make_run_dir(
|
||||
base,
|
||||
"20260309-RECENT",
|
||||
Some(serde_json::json!({
|
||||
"run_id": "recent-1",
|
||||
"workflow_name": "pipeline",
|
||||
"goal": "",
|
||||
"start_time": "2025-06-01T00:00:00Z",
|
||||
"node_count": 1,
|
||||
"edge_count": 0
|
||||
})),
|
||||
Some(serde_json::json!({
|
||||
"timestamp": &recent_ts,
|
||||
"status": "success",
|
||||
"duration_ms": 1000
|
||||
})),
|
||||
false,
|
||||
);
|
||||
|
||||
let args = RunsPruneArgs {
|
||||
filter: RunFilterArgs {
|
||||
before: Some("2026-01-01".into()),
|
||||
workflow: None,
|
||||
label: Vec::new(),
|
||||
orphans: false,
|
||||
},
|
||||
older_than: None,
|
||||
yes: true,
|
||||
};
|
||||
|
||||
prune_from(&args, base).unwrap();
|
||||
assert!(
|
||||
!dir.exists(),
|
||||
"explicit --before should bypass staleness filter"
|
||||
);
|
||||
}
|
||||
|
||||
// === Step 2: parse_duration tests ===
|
||||
|
||||
#[test]
|
||||
fn parse_duration_hours() {
|
||||
let dur = parse_duration("24h").unwrap();
|
||||
assert_eq!(dur, chrono::Duration::hours(24));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parse_duration_days() {
|
||||
let dur = parse_duration("7d").unwrap();
|
||||
assert_eq!(dur, chrono::Duration::days(7));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parse_duration_invalid() {
|
||||
assert!(parse_duration("").is_err());
|
||||
assert!(parse_duration("abc").is_err());
|
||||
}
|
||||
|
||||
// === Step 3: disk space reporting tests ===
|
||||
|
||||
#[test]
|
||||
fn format_size_display() {
|
||||
assert_eq!(format_size(0), "0 B");
|
||||
assert_eq!(format_size(512), "512 B");
|
||||
assert_eq!(format_size(1024), "1.0 KB");
|
||||
assert_eq!(format_size(1536), "1.5 KB");
|
||||
assert_eq!(format_size(1024 * 1024), "1.0 MB");
|
||||
assert_eq!(format_size(1024 * 1024 * 1024), "1.0 GB");
|
||||
assert_eq!(format_size(1024 * 1024 * 1024 + 512 * 1024 * 1024), "1.5 GB");
|
||||
}
|
||||
|
||||
// === Step 4: custom threshold test ===
|
||||
|
||||
#[test]
|
||||
fn prune_older_than_custom_threshold() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let base = tmp.path();
|
||||
|
||||
// Run completed 3 days ago — should be pruned with --older-than 2d
|
||||
let three_days_ago = (Utc::now() - chrono::Duration::days(3)).to_rfc3339();
|
||||
let dir_old = make_run_dir(
|
||||
base,
|
||||
"20260306-OLD",
|
||||
Some(serde_json::json!({
|
||||
"run_id": "old-1",
|
||||
"workflow_name": "pipeline",
|
||||
"goal": "",
|
||||
"start_time": &three_days_ago,
|
||||
"node_count": 1,
|
||||
"edge_count": 0
|
||||
})),
|
||||
Some(serde_json::json!({
|
||||
"timestamp": &three_days_ago,
|
||||
"status": "success",
|
||||
"duration_ms": 1000
|
||||
})),
|
||||
false,
|
||||
);
|
||||
|
||||
// Run completed 1 day ago — should NOT be pruned with --older-than 2d
|
||||
let one_day_ago = (Utc::now() - chrono::Duration::days(1)).to_rfc3339();
|
||||
let dir_recent = make_run_dir(
|
||||
base,
|
||||
"20260308-RECENT",
|
||||
Some(serde_json::json!({
|
||||
"run_id": "recent-1",
|
||||
"workflow_name": "pipeline",
|
||||
"goal": "",
|
||||
"start_time": &one_day_ago,
|
||||
"node_count": 1,
|
||||
"edge_count": 0
|
||||
})),
|
||||
Some(serde_json::json!({
|
||||
"timestamp": &one_day_ago,
|
||||
"status": "success",
|
||||
"duration_ms": 1000
|
||||
})),
|
||||
false,
|
||||
);
|
||||
|
||||
// Run completed 10 days ago — should be pruned with --older-than 7d
|
||||
let ten_days_ago = (Utc::now() - chrono::Duration::days(10)).to_rfc3339();
|
||||
let dir_very_old = make_run_dir(
|
||||
base,
|
||||
"20260227-VERYOLD",
|
||||
Some(serde_json::json!({
|
||||
"run_id": "very-old-1",
|
||||
"workflow_name": "pipeline",
|
||||
"goal": "",
|
||||
"start_time": &ten_days_ago,
|
||||
"node_count": 1,
|
||||
"edge_count": 0
|
||||
})),
|
||||
Some(serde_json::json!({
|
||||
"timestamp": &ten_days_ago,
|
||||
"status": "success",
|
||||
"duration_ms": 1000
|
||||
})),
|
||||
false,
|
||||
);
|
||||
|
||||
// With --older-than 7d, only the 10-day-old run should be pruned
|
||||
let args = RunsPruneArgs {
|
||||
filter: RunFilterArgs {
|
||||
before: None,
|
||||
workflow: None,
|
||||
label: Vec::new(),
|
||||
orphans: false,
|
||||
},
|
||||
older_than: Some(chrono::Duration::days(7)),
|
||||
yes: true,
|
||||
};
|
||||
|
||||
prune_from(&args, base).unwrap();
|
||||
assert!(dir_old.exists(), "3-day-old run should survive 7d threshold");
|
||||
assert!(dir_recent.exists(), "1-day-old run should survive 7d threshold");
|
||||
assert!(!dir_very_old.exists(), "10-day-old run should be pruned with 7d threshold");
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue