diff --git a/bin/box-ctl.py b/bin/box-ctl.py index 9951f47..f802f47 100755 --- a/bin/box-ctl.py +++ b/bin/box-ctl.py @@ -540,15 +540,56 @@ def _job_summary(name, path): } -def act_job_list(): +def act_job_list(include_archived=False): jobs = [] if JOBS_DIR.exists(): for p in sorted(JOBS_DIR.glob("*.json")): jobs.append(_job_summary(p.stem, p)) + if include_archived: + arch_dir = JOBS_DIR / "archive" + if arch_dir.exists(): + for p in sorted(arch_dir.glob("*.json")): + s = _job_summary(p.stem, p) + s["archived"] = True + jobs.append(s) audit("job-list") out(True, jobs=jobs) +def act_job_archive(name, force=False): + check_name(name) + script = BIN / "retention-archive-jobs.py" + cmd = [sys.executable, str(script), "archive", name, "--json"] + if force: + cmd.append("--force") + r = run(cmd) + try: + data = json.loads(r.stdout) + except Exception: + fail("ARCHIVE_FAILED", r.stderr or r.stdout) + if data.get("status") in ("archived", "dry_run"): + audit("job-archive", name) + out(True, **data) + else: + fail("ARCHIVE_FAILED", data.get("error") or data.get("status")) + + +def act_job_unarchive(name): + check_name(name) + script = BIN / "retention-archive-jobs.py" + cmd = [sys.executable, str(script), "unarchive", name, "--json"] + r = run(cmd) + try: + data = json.loads(r.stdout) + except Exception: + fail("UNARCHIVE_FAILED", r.stderr or r.stdout) + if data.get("status") in ("unarchived", "dry_run"): + audit("job-unarchive", name) + out(True, **data) + else: + fail("UNARCHIVE_FAILED", data.get("error") or data.get("status")) + + def act_job_get(name): check_name(name) p = JOBS_DIR / f"{name}.json" @@ -1821,6 +1862,176 @@ def act_ssh_info(name): out(True, **agent_md.get_ssh_info(name)) +SSH_CHECK_TIMEOUT = 25 + + +def build_ssh_sweep_script(): + """Return the python3 probe script executed on the jump host. + + Reads `kind:port` args (kind is `ssh` or `term`), connects to each + 127.0.0.1:port, grabs the SSH banner or terminal HTTP status line, + and prints one JSON object: {"ports": {port: {...}}}. + """ + return ( + "import json, socket, sys, time\n" + "out = {}\n" + "for arg in sys.argv[1:]:\n" + " try:\n" + " kind, port_s = arg.split(':', 1)\n" + " port = int(port_s)\n" + " except ValueError:\n" + " continue\n" + " rec = {'open': False, 'latency_ms': None, 'banner': '', 'http': ''}\n" + " t0 = time.time()\n" + " try:\n" + " s = socket.create_connection(('127.0.0.1', port), timeout=2.0)\n" + " except Exception:\n" + " out[str(port)] = rec\n" + " continue\n" + " rec['open'] = True\n" + " rec['latency_ms'] = int((time.time() - t0) * 1000)\n" + " try:\n" + " s.settimeout(2.0)\n" + " if kind == 'term':\n" + " s.sendall(b'GET / HTTP/1.0\\r\\n\\r\\n')\n" + " data = s.recv(128)\n" + " rec['http'] = data.split(b'\\n')[0].decode('utf-8', 'replace').strip()[:64]\n" + " else:\n" + " data = s.recv(128)\n" + " rec['banner'] = data.split(b'\\n')[0].decode('utf-8', 'replace').strip()[:64]\n" + " except Exception:\n" + " pass\n" + " try:\n" + " s.close()\n" + " except Exception:\n" + " pass\n" + " out[str(port)] = rec\n" + "print(json.dumps({'ports': out}))\n" + ) + + +def parse_ssh_sweep_result(sweep, tunnels): + """Map sweep {port: probe} output onto {account: health} via TUNNEL_PORTS. + + Unknown sweep ports are ignored; accounts with no sweep data keep + None (unknown) health. Always returns every account in tunnels. + """ + ports = sweep.get("ports", {}) if isinstance(sweep, dict) else {} + by_port = {} + for acct, info in tunnels.items(): + by_port[str(info.get("port"))] = (acct, "ssh") + by_port[str(info.get("terminal"))] = (acct, "term") + + accounts = {} + for acct, info in tunnels.items(): + accounts[acct] = { + "ssh_port": info.get("port"), + "ssh_up": None, + "ssh_latency_ms": None, + "ssh_banner": "", + "term_port": info.get("terminal"), + "term_up": None, + "term_latency_ms": None, + "term_http": "", + "container_user": info.get("user", "hatch"), + } + if not isinstance(ports, dict): + return accounts + for port_s, probe in ports.items(): + slot = by_port.get(str(port_s)) + if not slot or not isinstance(probe, dict): + continue + acct, kind = slot + ent = accounts.get(acct) + if ent is None: + continue + is_open = bool(probe.get("open")) + lat = probe.get("latency_ms") + try: + lat = int(lat) if lat is not None else None + except (TypeError, ValueError): + lat = None + if kind == "ssh": + ent["ssh_up"] = is_open + ent["ssh_latency_ms"] = lat if is_open else None + ent["ssh_banner"] = str(probe.get("banner") or "")[:64] + else: + ent["term_up"] = is_open + ent["term_latency_ms"] = lat if is_open else None + ent["term_http"] = str(probe.get("http") or "")[:64] + return accounts + + +def run_vm_port_sweep(jump_host, operator_user, tunnels, timeout=SSH_CHECK_TIMEOUT): + """Run one SSH to the jump host sweeping all tunnel ports. + + Returns (sweep_dict_or_None, error_str). sweep is the parsed + {"ports": {...}} payload; error is "" on success. + """ + sweep_args = [] + for _acct, info in tunnels.items(): + sweep_args.append(f"ssh:{info.get('port')}") + sweep_args.append(f"term:{info.get('terminal')}") + cmd = [ + "ssh", "-o", "BatchMode=yes", "-o", "ConnectTimeout=8", + "-o", "StrictHostKeyChecking=no", + ] + identity = os.environ.get("SSH_IDENTITY_FILE", "") + if identity: + cmd += ["-o", "IdentitiesOnly=yes", "-i", identity] + cmd += [ + f"{operator_user}@{jump_host}", + "python3", "-", + ] + sweep_args + try: + r = subprocess.run(cmd, input=build_ssh_sweep_script(), + capture_output=True, text=True, timeout=timeout) + except subprocess.TimeoutExpired: + return None, f"jump host {jump_host} sweep timed out after {timeout}s" + except FileNotFoundError: + return None, "local ssh binary not found" + except Exception as e: + return None, f"ssh to {jump_host} failed: {e}" + if r.returncode != 0: + err = (r.stderr or "").strip().splitlines() + hint = err[-1][:160] if err else f"exit {r.returncode}" + return None, f"jump host {jump_host} unreachable: {hint}" + try: + sweep = json.loads(r.stdout) + except Exception: + return None, f"jump host {jump_host} returned unparseable sweep output" + if not isinstance(sweep, dict) or "ports" not in sweep: + return None, f"jump host {jump_host} returned malformed sweep output" + return sweep, "" + + +def act_ssh_check(): + """Probe all container reverse-tunnel ports from the jump host. + + Single SSH connection, VM-side sweep. Always emits HTTP-200-style + ok:true with per-account health; jump failures surface as + jump_reachable:false (degraded-state data, not a fatal error). + """ + import time as _time + import agent_md + audit("ssh-check") + t0 = _time.time() + tunnels = dict(agent_md.TUNNEL_PORTS) + jump_host = os.environ.get("SSH_JUMP_HOST", "34.139.37.135") + operator_user = os.environ.get("OPERATOR_USER", "super") + sweep, err = run_vm_port_sweep(jump_host, operator_user, tunnels) + wall_ms = int((_time.time() - t0) * 1000) + accounts = parse_ssh_sweep_result(sweep or {}, tunnels) + if err: + out(True, jump_host=jump_host, operator_user=operator_user, + jump_reachable=False, error=err, accounts=accounts, + checked_at=utcnow(), latency_ms=wall_ms) + else: + out(True, jump_host=jump_host, operator_user=operator_user, + jump_reachable=True, accounts=accounts, + checked_at=utcnow(), latency_ms=wall_ms) + + MD_MAX_READ = 64 * 1024 MD_MAX_DIFF = 64 * 1024 MD_MAX_LIST = 200 @@ -2252,12 +2463,13 @@ onboard actions: onboard-connects consolidated active fleet and client onboard connects (read-only) (space-separated alias: onboard connects) -ssh actions (space-separated alias: ssh mint|list|show|ports|info): +ssh actions (space-separated alias: ssh mint|list|show|ports|info|check): ssh-mint [--force] mint a new SSH key ssh-list list minted keys ssh-show show key detail ssh-ports tunnel port inventory - ssh-info [name] connection coordinates""" + ssh-info [name] connection coordinates + ssh-check live tunnel health sweep via jump host""" def act_thread(op, agent, thread=None, title=None, limit=None, confirm=False): @@ -3723,7 +3935,17 @@ def main(argv): else: fail("BAD_NAME", "usage: timer list|status|create|delete|start|stop|enable|disable [...]") elif action == "job-list": - act_job_list() + include_archived = "--archived" in rest or "-a" in rest + act_job_list(include_archived=include_archived) + elif action == "job-archive": + if not rest or len(rest) > 2: + fail("BAD_NAME", "usage: job-archive [--force]") + force = "--force" in rest[1:] + act_job_archive(rest[0], force=force) + elif action == "job-unarchive": + if len(rest) != 1: + fail("BAD_NAME", "usage: job-unarchive ") + act_job_unarchive(rest[0]) elif action == "job-get": if len(rest) != 1: fail("BAD_NAME", "usage: job-get ") @@ -3861,7 +4083,7 @@ def main(argv): fail("BAD_NAME", "usage: loop-resolve [note]") note = rest[1] if len(rest) > 1 else None act_loop_resolve(rest[0], note=note) - elif action in ("ssh", "ssh-mint", "ssh-list", "ssh-show", "ssh-ports", "ssh-info"): + elif action in ("ssh", "ssh-mint", "ssh-list", "ssh-show", "ssh-ports", "ssh-info", "ssh-check"): if action == "ssh-mint": if not rest: fail("BAD_ARGS", "usage: ssh-mint [--force]") @@ -3876,9 +4098,11 @@ def main(argv): act_ssh_ports() elif action == "ssh-info": act_ssh_info(rest[0] if rest else None) + elif action == "ssh-check": + act_ssh_check() elif action == "ssh": - if not rest or rest[0] not in ("mint", "list", "show", "ports", "info"): - fail("BAD_NAME", "usage: ssh mint|list|show|ports|info [...]") + if not rest or rest[0] not in ("mint", "list", "show", "ports", "info", "check"): + fail("BAD_NAME", "usage: ssh mint|list|show|ports|info|check [...]") sub = rest[0] args = rest[1:] if sub == "mint": @@ -3895,6 +4119,8 @@ def main(argv): act_ssh_ports() elif sub == "info": act_ssh_info(args[0] if args else None) + elif sub == "check": + act_ssh_check() elif action in ("md", "md-audit", "md-list", "md-read", "md-write", "md-diff", "md-inject-drive", "md-sync-all", "md-amend", "md-append", "md-pull"): if action == "md-audit": diff --git a/bin/job-dispatch.py b/bin/job-dispatch.py index 5374fd5..374f603 100755 --- a/bin/job-dispatch.py +++ b/bin/job-dispatch.py @@ -103,6 +103,10 @@ def load_job(job_name): # Try .json first, then .yaml (for compatibility) job_file = JOBS_DIR / f"{job_name}.json" if not job_file.exists(): + archived_file = JOBS_DIR / "archive" / f"{job_name}.json" + if archived_file.exists(): + print(f"Error: Job '{job_name}' is archived at {archived_file}. Unarchive before dispatch (e.g. 'box job unarchive {job_name}').", file=sys.stderr) + sys.exit(1) job_file = JOBS_DIR / f"{job_name}.yaml" if not job_file.exists(): print(f"Error: Job '{job_name}' not found in {JOBS_DIR}", file=sys.stderr) diff --git a/bin/retention-archive-jobs.py b/bin/retention-archive-jobs.py new file mode 100755 index 0000000..bf2e99f --- /dev/null +++ b/bin/retention-archive-jobs.py @@ -0,0 +1,409 @@ +#!/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() diff --git a/bin/retention-run-rotations.sh b/bin/retention-run-rotations.sh index 4b7ad02..4017534 100755 --- a/bin/retention-run-rotations.sh +++ b/bin/retention-run-rotations.sh @@ -22,6 +22,7 @@ fail=0 "$BIN/retention-rotate-chat-history.sh" || fail=1 "$BIN/retention-rotate-job-log.sh" || fail=1 "$BIN/retention-archive-followups.py" || fail=1 +"$BIN/retention-archive-jobs.py" scan || fail=1 echo "--- verification ---" @@ -40,6 +41,7 @@ done echo "--- live sizes after run ---" ls -la "$ROOT/logs/chat-history.jsonl" "$ROOT/job-log.jsonl" "$ROOT/followups.json" 2>/dev/null || true +echo "active jobs: $(ls -1 "$ROOT/jobs"/*.json 2>/dev/null | wc -l), archived jobs: $(ls -1 "$ROOT/jobs/archive"/*.json 2>/dev/null | wc -l)" echo "=== retention-run done rc=$fail ===" exit $fail diff --git a/bin/super-cli.py b/bin/super-cli.py index 55cb6e4..095cf5a 100755 --- a/bin/super-cli.py +++ b/bin/super-cli.py @@ -960,7 +960,7 @@ def cmd_runtime(args): state = r["state"] if r["state"] == "approval-pending" and r["prompt_kind"]: state = "%s(%s/%s)" % (state, r["prompt_kind"], - r["prompt_key"]) + r["prompt_key"] or "…") if r["is_muse"]: approve = (badge_ok("YES") if r["auto_approve"] else badge_err("NO")) @@ -995,29 +995,30 @@ def cmd_runtime(args): print(c_red("Error: no such pane %s on %s" % (pane, sock)), file=sys.stderr) sys.exit(1) - t_args = ["send-keys", "-t", pane, keys] - if enter: - t_args.append("Enter") - r = mcw._tmux(sock, *t_args, timeout=10) - ok = r.returncode == 0 + # Verified send: keys + Enter in one tmux call arrive as a paste + # burst, which the muse composer takes as a newline instead of a + # submit. send_answer paces literal text and Enter apart (with a + # render check for composer-bound prompts) so the submit lands. + ok, detail = mcw.send_answer(sock, pane, keys, enter=enter, + kind=pre.get("prompt_kind"), sig=None) if as_json: print(json.dumps({ "ok": ok, "socket": sock, "pane": pane, "pre_state": pre["state"], "prompt_kind": pre["prompt_kind"], "sent": keys, "enter": enter, - "error": (r.stderr or r.stdout or "").strip() or None - if not ok else None}, indent=2)) + "verified": detail["verified"], + "retried": detail["retried"], + "error": None if ok else "tmux error"}, indent=2)) return print("\n Pane %s on %s [%s]" % ( c_bold(pane), sock, c_cyan(pre["state"]))) if ok: - print(" %s sent %r%s" % (c_green("✔"), - keys, " + Enter" if enter else "")) + how = " (verified)" if detail["verified"] else "" + print(" %s sent %r%s%s" % (c_green("✔"), keys, + " + Enter" if enter else "", how)) else: - print(" %s send failed: %s" % ( - c_red("✘"), - (r.stderr or r.stdout or "").strip() or "tmux error")) + print(" %s send failed: %s" % (c_red("✘"), "tmux error")) sys.exit(1) print() @@ -2845,6 +2846,8 @@ def cmd_thread_view(args): # --------------------------------------------------------------------------- def cmd_job_list(args): cmd = ["python3", str(BIN_DIR / "box-ctl.py"), "job-list"] + if getattr(args, "archived", False): + cmd.append("--archived") res = subprocess.run(cmd, capture_output=True, text=True) try: @@ -2858,23 +2861,63 @@ def cmd_job_list(args): return jobs = data.get("jobs", []) - print("\n" + c_bold(f"=== DEFINED JOBS ({len(jobs)}) ===") + "\n") + title_extra = " (INCLUDING ARCHIVED)" if getattr(args, "archived", False) else "" + print("\n" + c_bold(f"=== DEFINED JOBS ({len(jobs)}){title_extra} ===") + "\n") - headers = ["NAME", "AGENT", "SCHEDULE", "TIMEOUT", "DESCRIPTION"] + headers = ["NAME", "AGENT", "SCHEDULE", "TIMEOUT", "STATUS", "DESCRIPTION"] rows = [] for j in jobs: sched = j.get("schedule", "-") sched_disp = c_green(sched) if sched != "manual" else c_dim("manual") + is_arch = j.get("archived", False) + status_disp = c_yellow("ARCHIVED") if is_arch else c_green("ACTIVE") rows.append([ c_bold(j.get("name", "")), j.get("agent", "-"), sched_disp, f"{j.get('timeout', '-')}s", - (j.get("description") or "-")[:50] + status_disp, + (j.get("description") or "-")[:45] ]) print_table(headers, rows) - print("\n" + c_dim(" Commands: super job show | super job status | super job run | super job create ") + "\n") + print("\n" + c_dim(" Commands: box job show | box job archive | box job unarchive | box job run ") + "\n") + + +def cmd_job_archive(args): + cmd = ["python3", str(BIN_DIR / "box-ctl.py"), "job-archive", args.name] + if getattr(args, "force", False): + cmd.append("--force") + res = subprocess.run(cmd, capture_output=True, text=True) + try: + data = json.loads(res.stdout) + except Exception: + print(c_red("Failed to archive job:") + f"\n{res.stderr or res.stdout}") + return + if args.json: + print(json.dumps(data, indent=2)) + return + if data.get("ok"): + print(c_green(f"✔ Archived job '{args.name}' to jobs/archive/{args.name}.json")) + else: + print(c_red(f"Error archiving job '{args.name}': {data.get('error', 'unknown error')}")) + + +def cmd_job_unarchive(args): + cmd = ["python3", str(BIN_DIR / "box-ctl.py"), "job-unarchive", args.name] + res = subprocess.run(cmd, capture_output=True, text=True) + try: + data = json.loads(res.stdout) + except Exception: + print(c_red("Failed to unarchive job:") + f"\n{res.stderr or res.stdout}") + return + if args.json: + print(json.dumps(data, indent=2)) + return + if data.get("ok"): + print(c_green(f"✔ Unarchived job '{args.name}' to jobs/{args.name}.json")) + else: + print(c_red(f"Error unarchiving job '{args.name}': {data.get('error', 'unknown error')}")) def cmd_job_show(args): name = args.name @@ -4718,6 +4761,58 @@ def cmd_ssh_info(args): sys.stdout.write(res.stdout) +def cmd_ssh_check(args): + cmd = ["python3", str(BIN_DIR / "box-ctl.py"), "ssh-check"] + res = subprocess.run(cmd, capture_output=True, text=True) + if args.json: + sys.stdout.write(res.stdout) + return + try: + d = json.loads(res.stdout) + except Exception: + sys.stdout.write(res.stdout) + return + if not d.get("ok"): + print(c_red(f"✗ Failed: {d.get('error')}")) + return + jump = d.get("jump_host", "?") + reachable = d.get("jump_reachable") + wall = d.get("latency_ms") + print("\n" + c_bold(f"=== CONTAINER SSH TUNNEL HEALTH (JUMP: {jump}) ===") + "\n") + if not reachable: + print(c_red(f" ✗ Jump host unreachable: {d.get('error', 'unknown')}")) + print(c_dim(" Tunnel states below are UNKNOWN (last sweep failed).") + "\n") + elif wall is not None: + print(c_dim(f" Sweep completed in {wall}ms at {d.get('checked_at', '')}") + "\n") + headers = ["AGENT", "SSH", "LAT", "BANNER", "TERM", "STATE", "DIAL"] + rows = [] + accounts = d.get("accounts", {}) + + def _state(v): + if v is True: + return c_green("UP") + if v is False: + return c_red("DOWN") + return c_dim("?") + + for name, info in accounts.items(): + sport = info.get("ssh_port", "?") + tport = info.get("term_port", "?") + lat = info.get("ssh_latency_ms") + lat_s = f"{lat}ms" if lat is not None else "-" + banner = (info.get("ssh_banner") or "-")[:28] + therm = info.get("term_http") or "" + tstate = _state(info.get("term_up")) + if info.get("term_up") and therm: + tstate += c_dim(f" {therm[:18]}") + u = info.get("container_user", "hatch") + dial = f"ssh -J super@{jump} -p {sport} {u}@localhost" + rows.append([c_cyan(name), f":{sport} {_state(info.get('ssh_up'))}", + lat_s, banner, f":{tport}", tstate, c_yellow(dial)]) + print_table(headers, rows) + print() + + def cmd_md_audit(args): cmd = ["python3", str(BIN_DIR / "box-ctl.py"), "md-audit"] + (args.accounts or []) res = subprocess.run(cmd, capture_output=True, text=True) @@ -5931,6 +6026,7 @@ def build_parser(): job_sub = p_job.add_subparsers(dest="action") p_job_list = job_sub.add_parser("list", parents=[common], help="List defined jobs") + p_job_list.add_argument("--archived", "-a", action="store_true", help="Include archived jobs from jobs/archive/") p_job_show = job_sub.add_parser("show", parents=[common], help="Inspect full job definition and timer details") p_job_show.add_argument("name", help="Job name") @@ -5974,6 +6070,14 @@ def build_parser(): p_job_log.add_argument("-n", type=int, default=20, help="Number of events") p_job_log.add_argument("name", nargs="?", default=None, help="Filter by job name") + p_job_archive = job_sub.add_parser("archive", parents=[common], help="Archive a retired manual job to jobs/archive/") + p_job_archive.add_argument("name", help="Job name to archive") + p_job_archive.add_argument("--force", action="store_true", help="Force archive even if marked protected") + + p_job_unarchive = job_sub.add_parser("unarchive", parents=[common], help="Restore an archived job from jobs/archive/") + p_job_unarchive.add_argument("name", help="Job name to restore") + + # Domain: WEB p_web = subparsers.add_parser("web", parents=[common], help="Probe VM web surfaces (box.muse-dev.online), auth, sync") web_sub = p_web.add_subparsers(dest="action") @@ -6025,6 +6129,7 @@ def build_parser(): p_s_ports = ssh_sub.add_parser("ports", parents=[common], help="List container reverse tunnel port mappings") p_s_info = ssh_sub.add_parser("info", parents=[common], help="Show SSH dial-in commands and reverse tunnel coordinates") p_s_info.add_argument("name", help="Agent identity (e.g. 646, muse, pip)") + p_s_check = ssh_sub.add_parser("check", parents=[common], help="Live SSH/terminal tunnel health sweep via jump host") # Domain: MD (Hatch Agent Markdown & Operational Drive) p_md = subparsers.add_parser("md", parents=[common], help="Manage agent .md system files & inject operational DRIVE via Hatch") @@ -6480,6 +6585,10 @@ def main(): cmd_job_run(args) elif act == "log": cmd_job_log(args) + elif act == "archive": + cmd_job_archive(args) + elif act == "unarchive": + cmd_job_unarchive(args) else: parser.print_help() elif args.domain == "web": @@ -6520,6 +6629,8 @@ def main(): cmd_ssh_ports(args) elif act == "info": cmd_ssh_info(args) + elif act == "check": + cmd_ssh_check(args) else: p_ssh.print_help() elif args.domain == "md": diff --git a/docs/RETENTION-ROTATIONS-P2.md b/docs/RETENTION-ROTATIONS-P2.md new file mode 100644 index 0000000..cae56b1 --- /dev/null +++ b/docs/RETENTION-ROTATIONS-P2.md @@ -0,0 +1,70 @@ +# Retention Piece 2 — Job Definition Archival & Pruning + +Owner: operator-646 +Branch: `dev/operator-646/retention-rotations-p1` +Date: 2026-10-07 + +Extends the fleet retention framework from Piece 1 ([RETENTION-ROTATIONS-P1.md](file:///home/super/Projects/NetVM/docs/RETENTION-ROTATIONS-P1.md)) to prune and archive retired manual job definitions from `jobs/` into `jobs/archive/`. + +## Problem + +The root `jobs/` directory accumulated over 235 job definition files, primarily composed of retired `auto-work-*` exploratory batches and one-off test pipelines. Every invocation of `job-scheduler.py` (`glob.glob("jobs/*.json")`) and CLI tools was parsing all 235 files, causing unnecessary filesystem overhead and clutter. + +## Policy Implemented + +| Domain | Eligibility Criteria | Action | Target Layout | +|:---|:---|:---|:---| +| `jobs/*.json` | `schedule == "manual"` or empty, AND not protected | `git mv` archive | `jobs/archive/.json` (plain JSON, preserves history) | + +### Dynamic Safety Guardrails + +A job is **never** archived automatically if any of the following hold: +1. **Systemd User Units**: Registered as an `ExecStart` target in `~/.config/systemd/user/job-*.service` (e.g. `work-finder`, `heartbeat`, `muse-auditor`, `opm-swarm-harvest`, `autonomy-pulse-*`). +2. **Active Pipelines**: Associated with a running, dispatched, or pending run in `pipelines.json` (e.g. `ops-audit-*`, `pipe-demo-*`). +3. **Chained Dependencies**: Linked via `chain_next`, `on_success`, or `on_failure` from any active protected job. +4. **Core Fleet Baseline**: Hardcoded protective baseline: `heartbeat`, `refine-system`, `canary-test`. + +## Tooling & Architecture + +- `bin/retention-archive-jobs.py`: Standalone CLI driver supporting `scan`, `archive`, `unarchive`, and `list` with `--dry-run`, `--force`, and `--json`. +- `bin/retention-run-rotations.sh`: Hourly rotation driver now runs `retention-archive-jobs.py scan` as the 4th phase of the retention pipeline. +- `bin/box-ctl.py`: + - `job-list`: Excludes archived jobs by default; includes them when passed `--archived` / `-a`. + - `job-archive `: Archives a job via `retention-archive-jobs.py`. + - `job-unarchive `: Restores an archived job to active `jobs/`. +- `bin/super-cli.py`: Exposes `box job archive `, `box job unarchive `, and `box job list --archived`. +- `bin/job-dispatch.py`: Guards against direct dispatch of archived jobs; fails fast with an explicit unarchive prompt. + +## Operations & Verification + +Run on-demand scan: +```bash +bin/retention-archive-jobs.py scan [--dry-run] +``` + +Manual archive / unarchive: +```bash +box job archive +box job unarchive +``` + +List archived jobs: +```bash +box job list --archived +``` + +Execute full retention cycle: +```bash +bin/retention-run-rotations.sh +``` + +## First Run Results (2026-10-07) + +- Total jobs scanned: 235 +- Dynamically protected jobs: 13 +- Eligible retired jobs identified: 17 +- Jobs archived to `jobs/archive/`: + `646-opm-watch`, `646-pip-sync`, `646-sidechat-task`, `auto-work-queue-f02`, `auto-work-queue-f06`, `auto-work-queue-f10`, `auto-work-queue-f14`, `auto-work-queue-f18`, `auto-work-xop-e02`, `auto-work-xop-e06`, `auto-work-xop-e10`, `auto-work-xop-e14`, `auto-work-xop-e18`, `mainloop-p1-pilot`, `mainloop-p2-noswitcher`, `mainloop-p3-bridge`, `mainloop-p4-steady`. +- Committed in git as: `fff5556` (`chore(retention): archive retired jobs [...]`). +- Active jobs remaining: 218. +- Unit tests: 10 tests in `tests/test_retention_archive_jobs.py` (all green). diff --git a/tests/test_retention_archive_jobs.py b/tests/test_retention_archive_jobs.py new file mode 100644 index 0000000..56852a0 --- /dev/null +++ b/tests/test_retention_archive_jobs.py @@ -0,0 +1,157 @@ +#!/usr/bin/env python3 +"""test_retention_archive_jobs.py — Unit tests for Retention Piece 2: Job Archival & Pruning.""" + +import json +import os +import shutil +import sys +import tempfile +import unittest +from pathlib import Path +from unittest.mock import patch + +REPO_ROOT = Path(__file__).resolve().parent.parent +sys.path.insert(0, str(REPO_ROOT / "bin")) + +import importlib.util +spec = importlib.util.spec_from_file_location("retention_archive_jobs", str(REPO_ROOT / "bin" / "retention-archive-jobs.py")) +raj = importlib.util.module_from_spec(spec) +spec.loader.exec_module(raj) + + +class TestRetentionArchiveJobs(unittest.TestCase): + def setUp(self): + self.test_dir = tempfile.mkdtemp() + self.root = Path(self.test_dir) + self.jobs_dir = self.root / "jobs" + self.archive_dir = self.jobs_dir / "archive" + self.jobs_dir.mkdir(parents=True) + + def tearDown(self): + shutil.rmtree(self.test_dir, ignore_errors=True) + + def _create_job(self, name, schedule="manual", extra=None): + data = { + "name": name, + "description": f"Test job {name}", + "agent": "pip", + "schedule": schedule, + "timeout": 300, + } + if extra: + data.update(extra) + p = self.jobs_dir / f"{name}.json" + p.write_text(json.dumps(data, indent=2)) + return p + + def test_core_protected_jobs(self): + protected = raj.get_protected_jobs(self.root) + self.assertIn("heartbeat", protected) + self.assertIn("refine-system", protected) + self.assertIn("canary-test", protected) + + def test_systemd_unit_protection(self): + sdir = self.root / "systemd" + sdir.mkdir() + svc = sdir / "job-custom-worker.service" + svc.write_text("[Service]\nExecStart=/home/super/Projects/NetVM/bin/job-dispatch.py custom-worker\n") + + protected = raj.get_protected_jobs(self.root) + self.assertIn("custom-worker", protected) + self.assertIn("systemd unit", protected["custom-worker"]) + + def test_active_pipeline_protection(self): + pipe_file = self.root / "pipelines.json" + pipe_data = { + "run-001": { + "status": "running", + "pipeline_name": "live-pipe-root", + "steps": [{"job_name": "live-pipe-step2"}] + }, + "run-002": { + "status": "completed", + "pipeline_name": "dead-pipe-root", + "steps": [{"job_name": "dead-pipe-step"}] + } + } + pipe_file.write_text(json.dumps(pipe_data)) + + protected = raj.get_protected_jobs(self.root) + self.assertIn("live-pipe-root", protected) + self.assertIn("live-pipe-step2", protected) + self.assertNotIn("dead-pipe-root", protected) + + def test_chain_target_protection(self): + self._create_job("canary-test", schedule="manual", extra={"on_success": "canary-followup"}) + self._create_job("canary-followup", schedule="manual") + + protected = raj.get_protected_jobs(self.root) + self.assertIn("canary-followup", protected) + self.assertIn("linked by canary-test", protected["canary-followup"]) + + def test_eligibility_criteria(self): + manual_job = self._create_job("retired-task", schedule="manual") + cron_job = self._create_job("hourly-task", schedule="*/5 * * * *") + protected_job = self._create_job("heartbeat", schedule="manual") + + prot = {"heartbeat": "core baseline"} + + el1, _ = raj.is_eligible(manual_job, prot) + self.assertTrue(el1) + + el2, _ = raj.is_eligible(cron_job, prot) + self.assertFalse(el2) + + el3, _ = raj.is_eligible(protected_job, prot) + self.assertFalse(el3) + + def test_archive_and_unarchive_roundtrip(self): + self._create_job("auto-task-x1", schedule="manual") + + with patch.object(raj, "is_git_tracked", return_value=False): + # Archive + res_arch = raj.archive_job("auto-task-x1", dry_run=False, root_dir=self.root) + self.assertEqual(res_arch["status"], "archived") + self.assertFalse((self.jobs_dir / "auto-task-x1.json").exists()) + self.assertTrue((self.archive_dir / "auto-task-x1.json").exists()) + + # Unarchive + res_unarch = raj.unarchive_job("auto-task-x1", dry_run=False, root_dir=self.root) + self.assertEqual(res_unarch["status"], "unarchived") + self.assertTrue((self.jobs_dir / "auto-task-x1.json").exists()) + self.assertFalse((self.archive_dir / "auto-task-x1.json").exists()) + + def test_refuse_archive_protected_job(self): + self._create_job("heartbeat", schedule="manual") + res = raj.archive_job("heartbeat", dry_run=False, root_dir=self.root) + self.assertEqual(res["status"], "rejected") + self.assertTrue((self.jobs_dir / "heartbeat.json").exists()) + + def test_force_archive_protected_job(self): + self._create_job("heartbeat", schedule="manual") + with patch.object(raj, "is_git_tracked", return_value=False): + res = raj.archive_job("heartbeat", dry_run=False, force=True, root_dir=self.root) + self.assertEqual(res["status"], "archived") + self.assertTrue((self.archive_dir / "heartbeat.json").exists()) + + def test_scan_and_archive_dry_run(self): + self._create_job("task-1", schedule="manual") + self._create_job("task-2", schedule="*/10 * * * *") + + res = raj.scan_and_archive(dry_run=True, root_dir=self.root) + self.assertEqual(res["eligible"], 1) + self.assertEqual(res["archived"], 1) + self.assertTrue((self.jobs_dir / "task-1.json").exists()) + self.assertFalse((self.archive_dir / "task-1.json").exists()) + + def test_list_archived(self): + self._create_job("archived-task", schedule="manual") + with patch.object(raj, "is_git_tracked", return_value=False): + raj.archive_job("archived-task", dry_run=False, root_dir=self.root) + archived = raj.list_archived(root_dir=self.root) + self.assertEqual(len(archived), 1) + self.assertEqual(archived[0]["name"], "archived-task") + + +if __name__ == "__main__": + unittest.main()