Files
box/bin/retention-archive-jobs.py
T

410 lines
16 KiB
Python
Raw Normal View History

#!/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()