410 lines
16 KiB
Python
410 lines
16 KiB
Python
|
|
#!/usr/bin/env python3
|
||
|
|
"""retention-archive-jobs.py — Retention Piece 2: Job definition archival & pruning.
|
||
|
|
|
||
|
|
Archives retired manual job definitions from jobs/ to jobs/archive/ via git mv,
|
||
|
|
preventing scheduler parsing overhead while preserving full git history and readability.
|
||
|
|
|
||
|
|
Guards active jobs dynamically:
|
||
|
|
1. Systemd user units (~/.config/systemd/user/job-*.service)
|
||
|
|
2. Active pipeline runs in pipelines.json (and linked on_success/on_failure/chain_next)
|
||
|
|
3. Core baseline fleet jobs (heartbeat, refine-system, canary-test)
|
||
|
|
|
||
|
|
Usage:
|
||
|
|
retention-archive-jobs.py scan [--dry-run] [--limit N] [--json]
|
||
|
|
retention-archive-jobs.py archive <name> [--dry-run] [--force] [--json]
|
||
|
|
retention-archive-jobs.py unarchive <name> [--dry-run] [--json]
|
||
|
|
retention-archive-jobs.py list [--json]
|
||
|
|
"""
|
||
|
|
|
||
|
|
import argparse
|
||
|
|
import json
|
||
|
|
import os
|
||
|
|
import re
|
||
|
|
import subprocess
|
||
|
|
import sys
|
||
|
|
from pathlib import Path
|
||
|
|
|
||
|
|
NETVM_ROOT = Path(os.environ.get("NETVM_ROOT", "/home/super/Projects/NetVM"))
|
||
|
|
JOBS_DIR = NETVM_ROOT / "jobs"
|
||
|
|
ARCHIVE_DIR = JOBS_DIR / "archive"
|
||
|
|
PIPELINES_FILE = NETVM_ROOT / "pipelines.json"
|
||
|
|
SYSTEMD_USER_DIR = Path(os.path.expanduser("~/.config/systemd/user"))
|
||
|
|
|
||
|
|
CORE_PROTECTED = frozenset({"heartbeat", "refine-system", "canary-test"})
|
||
|
|
|
||
|
|
|
||
|
|
def run_git(*args, cwd=None):
|
||
|
|
"""Run git command safely and return (returncode, stdout, stderr)."""
|
||
|
|
cwd = cwd or NETVM_ROOT
|
||
|
|
try:
|
||
|
|
r = subprocess.run(["git"] + list(args), cwd=cwd, capture_output=True, text=True, check=False)
|
||
|
|
return r.returncode, r.stdout.strip(), r.stderr.strip()
|
||
|
|
except Exception as e:
|
||
|
|
return 1, "", str(e)
|
||
|
|
|
||
|
|
|
||
|
|
def is_git_tracked(rel_path, root_dir=None):
|
||
|
|
"""Check if a file is tracked in git."""
|
||
|
|
rc, stdout, _ = run_git("ls-files", str(rel_path), cwd=root_dir)
|
||
|
|
return rc == 0 and bool(stdout.strip())
|
||
|
|
|
||
|
|
|
||
|
|
def get_protected_jobs(root_dir=None):
|
||
|
|
"""Return dict of {job_name: reason} for all dynamically protected jobs."""
|
||
|
|
root = Path(root_dir) if root_dir else NETVM_ROOT
|
||
|
|
jobs_dir = root / "jobs"
|
||
|
|
protected = {k: "core baseline" for k in CORE_PROTECTED}
|
||
|
|
|
||
|
|
# 1. Systemd user timers / services
|
||
|
|
systemd_dirs = [SYSTEMD_USER_DIR, root / "systemd"]
|
||
|
|
for sdir in systemd_dirs:
|
||
|
|
if not sdir.exists():
|
||
|
|
continue
|
||
|
|
for sfile in sdir.glob("job-*.service"):
|
||
|
|
try:
|
||
|
|
content = sfile.read_text()
|
||
|
|
for line in content.splitlines():
|
||
|
|
if "job-dispatch.py" in line and line.strip().startswith("ExecStart="):
|
||
|
|
m = re.search(r"job-dispatch\.py\s+([A-Za-z0-9_-]+)", line)
|
||
|
|
if m:
|
||
|
|
protected[m.group(1)] = f"systemd unit: {sfile.name}"
|
||
|
|
except Exception:
|
||
|
|
pass
|
||
|
|
|
||
|
|
# 2. In-flight pipelines from pipelines.json
|
||
|
|
pipe_file = root / "pipelines.json"
|
||
|
|
if pipe_file.exists():
|
||
|
|
try:
|
||
|
|
data = json.loads(pipe_file.read_text())
|
||
|
|
for run_id, run_data in data.items():
|
||
|
|
status = str(run_data.get("status", "")).lower()
|
||
|
|
if status in ("running", "dispatched", "pending"):
|
||
|
|
pname = run_data.get("pipeline_name")
|
||
|
|
if pname:
|
||
|
|
protected[pname] = f"active pipeline {run_id}"
|
||
|
|
for step in run_data.get("steps", []):
|
||
|
|
sjob = step.get("job_name")
|
||
|
|
if sjob:
|
||
|
|
protected[sjob] = f"active pipeline {run_id} (step)"
|
||
|
|
except Exception:
|
||
|
|
pass
|
||
|
|
|
||
|
|
# 3. Chained job targets (transitive expansion of on_success, on_failure, chain_next)
|
||
|
|
to_expand = list(protected.keys())
|
||
|
|
visited = set()
|
||
|
|
while to_expand:
|
||
|
|
cur = to_expand.pop(0)
|
||
|
|
if cur in visited:
|
||
|
|
continue
|
||
|
|
visited.add(cur)
|
||
|
|
jpath = jobs_dir / f"{cur}.json"
|
||
|
|
if not jpath.exists():
|
||
|
|
continue
|
||
|
|
try:
|
||
|
|
d = json.loads(jpath.read_text())
|
||
|
|
for field in ("chain_next", "on_success", "on_failure"):
|
||
|
|
target = d.get(field)
|
||
|
|
if isinstance(target, str) and target.strip():
|
||
|
|
tname = target.strip()
|
||
|
|
# Skip common literal alerts like 'alert' unless an actual job definition exists
|
||
|
|
if tname == "alert" and not (jobs_dir / "alert.json").exists():
|
||
|
|
continue
|
||
|
|
if tname not in protected:
|
||
|
|
protected[tname] = f"linked by {cur} ({field})"
|
||
|
|
if tname not in visited:
|
||
|
|
to_expand.append(tname)
|
||
|
|
except Exception:
|
||
|
|
pass
|
||
|
|
|
||
|
|
return protected
|
||
|
|
|
||
|
|
|
||
|
|
def is_eligible(job_path, protected_jobs):
|
||
|
|
"""Determine if a job file is eligible for archival. Returns (bool, reason)."""
|
||
|
|
try:
|
||
|
|
data = json.loads(job_path.read_text())
|
||
|
|
except Exception as e:
|
||
|
|
return False, f"unreadable/invalid JSON: {e}"
|
||
|
|
|
||
|
|
name = data.get("name") or job_path.stem
|
||
|
|
if name in protected_jobs:
|
||
|
|
return False, f"protected: {protected_jobs[name]}"
|
||
|
|
|
||
|
|
schedule = str(data.get("schedule", "")).strip().lower()
|
||
|
|
if schedule == "manual" or not schedule:
|
||
|
|
return True, "manual schedule and not protected"
|
||
|
|
|
||
|
|
return False, f"active recurring schedule: {schedule}"
|
||
|
|
|
||
|
|
|
||
|
|
def archive_job(name, dry_run=False, force=False, root_dir=None):
|
||
|
|
"""Archive a single job by name. Returns dict with operation details."""
|
||
|
|
root = Path(root_dir) if root_dir else NETVM_ROOT
|
||
|
|
jobs_dir = root / "jobs"
|
||
|
|
archive_dir = jobs_dir / "archive"
|
||
|
|
src_file = jobs_dir / f"{name}.json"
|
||
|
|
dest_file = archive_dir / f"{name}.json"
|
||
|
|
|
||
|
|
if not src_file.exists():
|
||
|
|
if dest_file.exists():
|
||
|
|
return {"name": name, "status": "already_archived", "path": str(dest_file)}
|
||
|
|
return {"name": name, "status": "error", "error": f"Job file not found: {src_file}"}
|
||
|
|
|
||
|
|
# Verify JSON syntax
|
||
|
|
try:
|
||
|
|
content = json.loads(src_file.read_text())
|
||
|
|
except Exception as e:
|
||
|
|
return {"name": name, "status": "error", "error": f"Invalid source JSON: {e}"}
|
||
|
|
|
||
|
|
protected_jobs = get_protected_jobs(root)
|
||
|
|
if name in protected_jobs and not force:
|
||
|
|
return {"name": name, "status": "rejected", "error": f"Job is protected: {protected_jobs[name]}"}
|
||
|
|
|
||
|
|
if dry_run:
|
||
|
|
return {
|
||
|
|
"name": name,
|
||
|
|
"status": "dry_run",
|
||
|
|
"action": "would_archive",
|
||
|
|
"from": str(src_file),
|
||
|
|
"to": str(dest_file),
|
||
|
|
}
|
||
|
|
|
||
|
|
archive_dir.mkdir(parents=True, exist_ok=True)
|
||
|
|
tracked = is_git_tracked(src_file.relative_to(root), root_dir=root)
|
||
|
|
|
||
|
|
if tracked:
|
||
|
|
rc, _, err = run_git("mv", str(src_file), str(dest_file), cwd=root)
|
||
|
|
if rc != 0:
|
||
|
|
# Fallback to direct move + git add/rm
|
||
|
|
src_file.rename(dest_file)
|
||
|
|
run_git("add", str(dest_file), cwd=root)
|
||
|
|
run_git("rm", "-q", str(src_file), cwd=root)
|
||
|
|
else:
|
||
|
|
src_file.rename(dest_file)
|
||
|
|
|
||
|
|
# Verify destination file
|
||
|
|
try:
|
||
|
|
json.loads(dest_file.read_text())
|
||
|
|
except Exception as e:
|
||
|
|
return {"name": name, "status": "error", "error": f"Destination JSON verification failed: {e}"}
|
||
|
|
|
||
|
|
return {
|
||
|
|
"name": name,
|
||
|
|
"status": "archived",
|
||
|
|
"from": str(src_file),
|
||
|
|
"to": str(dest_file),
|
||
|
|
"git_tracked": tracked,
|
||
|
|
}
|
||
|
|
|
||
|
|
|
||
|
|
def unarchive_job(name, dry_run=False, root_dir=None):
|
||
|
|
"""Restore an archived job back to jobs/. Returns dict with operation details."""
|
||
|
|
root = Path(root_dir) if root_dir else NETVM_ROOT
|
||
|
|
jobs_dir = root / "jobs"
|
||
|
|
archive_dir = jobs_dir / "archive"
|
||
|
|
src_file = archive_dir / f"{name}.json"
|
||
|
|
dest_file = jobs_dir / f"{name}.json"
|
||
|
|
|
||
|
|
if not src_file.exists():
|
||
|
|
if dest_file.exists():
|
||
|
|
return {"name": name, "status": "already_active", "path": str(dest_file)}
|
||
|
|
return {"name": name, "status": "error", "error": f"Archived job not found: {src_file}"}
|
||
|
|
|
||
|
|
try:
|
||
|
|
content = json.loads(src_file.read_text())
|
||
|
|
except Exception as e:
|
||
|
|
return {"name": name, "status": "error", "error": f"Invalid archived JSON: {e}"}
|
||
|
|
|
||
|
|
if dry_run:
|
||
|
|
return {
|
||
|
|
"name": name,
|
||
|
|
"status": "dry_run",
|
||
|
|
"action": "would_unarchive",
|
||
|
|
"from": str(src_file),
|
||
|
|
"to": str(dest_file),
|
||
|
|
}
|
||
|
|
|
||
|
|
tracked = is_git_tracked(src_file.relative_to(root), root_dir=root)
|
||
|
|
|
||
|
|
if tracked:
|
||
|
|
rc, _, _ = run_git("mv", str(src_file), str(dest_file), cwd=root)
|
||
|
|
if rc != 0:
|
||
|
|
src_file.rename(dest_file)
|
||
|
|
run_git("add", str(dest_file), cwd=root)
|
||
|
|
run_git("rm", "-q", str(src_file), cwd=root)
|
||
|
|
else:
|
||
|
|
src_file.rename(dest_file)
|
||
|
|
|
||
|
|
try:
|
||
|
|
json.loads(dest_file.read_text())
|
||
|
|
except Exception as e:
|
||
|
|
return {"name": name, "status": "error", "error": f"Destination JSON verification failed: {e}"}
|
||
|
|
|
||
|
|
return {
|
||
|
|
"name": name,
|
||
|
|
"status": "unarchived",
|
||
|
|
"from": str(src_file),
|
||
|
|
"to": str(dest_file),
|
||
|
|
"git_tracked": tracked,
|
||
|
|
}
|
||
|
|
|
||
|
|
|
||
|
|
def scan_and_archive(dry_run=False, limit=None, root_dir=None):
|
||
|
|
"""Scan jobs/ for all eligible jobs and archive them."""
|
||
|
|
root = Path(root_dir) if root_dir else NETVM_ROOT
|
||
|
|
jobs_dir = root / "jobs"
|
||
|
|
archive_dir = jobs_dir / "archive"
|
||
|
|
|
||
|
|
protected_jobs = get_protected_jobs(root)
|
||
|
|
eligible = []
|
||
|
|
skipped = []
|
||
|
|
|
||
|
|
for jpath in sorted(jobs_dir.glob("*.json")):
|
||
|
|
name = jpath.stem
|
||
|
|
is_el, reason = is_eligible(jpath, protected_jobs)
|
||
|
|
if is_el:
|
||
|
|
eligible.append((name, jpath))
|
||
|
|
else:
|
||
|
|
skipped.append((name, reason))
|
||
|
|
|
||
|
|
to_process = eligible[:limit] if limit else eligible
|
||
|
|
results = []
|
||
|
|
|
||
|
|
for name, _ in to_process:
|
||
|
|
res = archive_job(name, dry_run=dry_run, root_dir=root)
|
||
|
|
results.append(res)
|
||
|
|
|
||
|
|
archived_count = sum(1 for r in results if r["status"] in ("archived", "dry_run"))
|
||
|
|
|
||
|
|
# Commit git changes if not dry_run and we actually archived tracked files
|
||
|
|
commit_sha = None
|
||
|
|
if not dry_run and archived_count > 0:
|
||
|
|
names_str = ", ".join(r["name"] for r in results if r["status"] == "archived")
|
||
|
|
commit_msg = f"chore(retention): archive retired jobs [{names_str}]"
|
||
|
|
rc, out, err = run_git("commit", "-m", commit_msg, cwd=root)
|
||
|
|
if rc == 0:
|
||
|
|
_, sha, _ = run_git("rev-parse", "--short", "HEAD", cwd=root)
|
||
|
|
commit_sha = sha
|
||
|
|
|
||
|
|
return {
|
||
|
|
"scanned": len(list(jobs_dir.glob("*.json"))),
|
||
|
|
"eligible": len(eligible),
|
||
|
|
"archived": archived_count,
|
||
|
|
"protected_total": len(protected_jobs),
|
||
|
|
"dry_run": dry_run,
|
||
|
|
"commit": commit_sha,
|
||
|
|
"results": results,
|
||
|
|
}
|
||
|
|
|
||
|
|
|
||
|
|
def list_archived(root_dir=None):
|
||
|
|
"""List all currently archived jobs."""
|
||
|
|
root = Path(root_dir) if root_dir else NETVM_ROOT
|
||
|
|
archive_dir = root / "jobs" / "archive"
|
||
|
|
if not archive_dir.exists():
|
||
|
|
return []
|
||
|
|
|
||
|
|
jobs = []
|
||
|
|
for p in sorted(archive_dir.glob("*.json")):
|
||
|
|
try:
|
||
|
|
d = json.loads(p.read_text())
|
||
|
|
jobs.append({
|
||
|
|
"name": d.get("name") or p.stem,
|
||
|
|
"agent": d.get("agent", "-"),
|
||
|
|
"schedule": d.get("schedule", "-"),
|
||
|
|
"description": d.get("description", ""),
|
||
|
|
"archived_path": str(p),
|
||
|
|
})
|
||
|
|
except Exception:
|
||
|
|
jobs.append({"name": p.stem, "error": "unreadable JSON", "archived_path": str(p)})
|
||
|
|
return jobs
|
||
|
|
|
||
|
|
|
||
|
|
def main():
|
||
|
|
parser = argparse.ArgumentParser(description="NetVM Job Definition Retention & Archival Driver (P2)")
|
||
|
|
sub = parser.add_subparsers(dest="subcommand", required=True)
|
||
|
|
|
||
|
|
# scan
|
||
|
|
p_scan = sub.add_parser("scan", help="Scan jobs/ and archive all eligible retired manual jobs")
|
||
|
|
p_scan.add_argument("--dry-run", action="store_true", help="Print actions without modifying files")
|
||
|
|
p_scan.add_argument("--limit", type=int, default=None, help="Max jobs to archive in this run")
|
||
|
|
p_scan.add_argument("--json", action="store_true", help="Output machine-readable JSON")
|
||
|
|
|
||
|
|
# archive
|
||
|
|
p_arch = sub.add_parser("archive", help="Archive a specific job by name")
|
||
|
|
p_arch.add_argument("name", help="Job name")
|
||
|
|
p_arch.add_argument("--dry-run", action="store_true", help="Simulate without modifying files")
|
||
|
|
p_arch.add_argument("--force", action="store_true", help="Force archive even if marked protected")
|
||
|
|
p_arch.add_argument("--commit", action="store_true", help="Create a git commit for the archive move")
|
||
|
|
p_arch.add_argument("--json", action="store_true", help="Output machine-readable JSON")
|
||
|
|
|
||
|
|
# unarchive
|
||
|
|
p_unarch = sub.add_parser("unarchive", help="Restore an archived job to jobs/")
|
||
|
|
p_unarch.add_argument("name", help="Job name")
|
||
|
|
p_unarch.add_argument("--dry-run", action="store_true", help="Simulate without modifying files")
|
||
|
|
p_unarch.add_argument("--commit", action="store_true", help="Create a git commit for the unarchive move")
|
||
|
|
p_unarch.add_argument("--json", action="store_true", help="Output machine-readable JSON")
|
||
|
|
|
||
|
|
# list
|
||
|
|
p_list = sub.add_parser("list", help="List archived jobs in jobs/archive/")
|
||
|
|
p_list.add_argument("--json", action="store_true", help="Output machine-readable JSON")
|
||
|
|
|
||
|
|
args = parser.parse_args()
|
||
|
|
|
||
|
|
if args.subcommand == "scan":
|
||
|
|
res = scan_and_archive(dry_run=args.dry_run, limit=args.limit)
|
||
|
|
if args.json:
|
||
|
|
print(json.dumps(res, indent=2))
|
||
|
|
else:
|
||
|
|
mode = " [DRY-RUN]" if args.dry_run else ""
|
||
|
|
print(f"=== Retention Job Archival Scan{mode} ===")
|
||
|
|
print(f"Scanned jobs: {res['scanned']}")
|
||
|
|
print(f"Eligible for archival: {res['eligible']}")
|
||
|
|
print(f"Archived count: {res['archived']}")
|
||
|
|
if res.get("commit"):
|
||
|
|
print(f"Committed as: {res['commit']}")
|
||
|
|
for item in res["results"]:
|
||
|
|
status = item["status"]
|
||
|
|
print(f" • {item['name']:25} -> {status}")
|
||
|
|
print("========================================")
|
||
|
|
|
||
|
|
elif args.subcommand == "archive":
|
||
|
|
res = archive_job(args.name, dry_run=args.dry_run, force=args.force)
|
||
|
|
if not args.dry_run and args.commit and res.get("status") == "archived":
|
||
|
|
run_git("commit", "-m", f"chore(retention): archive job {args.name}")
|
||
|
|
if args.json:
|
||
|
|
print(json.dumps(res, indent=2))
|
||
|
|
else:
|
||
|
|
if res.get("status") in ("archived", "dry_run"):
|
||
|
|
print(f"✔ Job '{args.name}' archived to jobs/archive/{args.name}.json")
|
||
|
|
else:
|
||
|
|
print(f"✖ Failed to archive job '{args.name}': {res.get('error') or res.get('status')}", file=sys.stderr)
|
||
|
|
sys.exit(1)
|
||
|
|
|
||
|
|
elif args.subcommand == "unarchive":
|
||
|
|
res = unarchive_job(args.name, dry_run=args.dry_run)
|
||
|
|
if not args.dry_run and args.commit and res.get("status") == "unarchived":
|
||
|
|
run_git("commit", "-m", f"chore(retention): unarchive job {args.name}")
|
||
|
|
if args.json:
|
||
|
|
print(json.dumps(res, indent=2))
|
||
|
|
else:
|
||
|
|
if res.get("status") in ("unarchived", "dry_run"):
|
||
|
|
print(f"✔ Job '{args.name}' unarchived to jobs/{args.name}.json")
|
||
|
|
else:
|
||
|
|
print(f"✖ Failed to unarchive job '{args.name}': {res.get('error') or res.get('status')}", file=sys.stderr)
|
||
|
|
sys.exit(1)
|
||
|
|
|
||
|
|
elif args.subcommand == "list":
|
||
|
|
items = list_archived()
|
||
|
|
if args.json:
|
||
|
|
print(json.dumps(items, indent=2))
|
||
|
|
else:
|
||
|
|
print(f"\n=== ARCHIVED JOBS ({len(items)}) ===")
|
||
|
|
for item in items:
|
||
|
|
print(f" • {item['name']:25} [{item.get('agent', '-')}] {item.get('description', '')[:50]}")
|
||
|
|
print()
|
||
|
|
|
||
|
|
|
||
|
|
if __name__ == "__main__":
|
||
|
|
main()
|