#!/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 [--dry-run] [--force] [--json] retention-archive-jobs.py unarchive [--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()