feat(retention): implement Piece 2 job archival, CLI wiring, and rotation driver

This commit is contained in:
operator
2026-10-07 03:28:38 +00:00
parent fff5556eb6
commit c11d1d83ae
7 changed files with 1003 additions and 24 deletions
+233 -7
View File
@@ -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 <name> [--force] mint a new SSH key
ssh-list list minted keys
ssh-show <name> 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 <name> [--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 <name>")
act_job_unarchive(rest[0])
elif action == "job-get":
if len(rest) != 1:
fail("BAD_NAME", "usage: job-get <name>")
@@ -3861,7 +4083,7 @@ def main(argv):
fail("BAD_NAME", "usage: loop-resolve <dm_id> [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 <name> [--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":
+4
View File
@@ -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)
+409
View File
@@ -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 <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()
+2
View File
@@ -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
+128 -17
View File
@@ -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 <name> | super job status <name> | super job run <name> | super job create <name>") + "\n")
print("\n" + c_dim(" Commands: box job show <name> | box job archive <name> | box job unarchive <name> | box job run <name>") + "\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":
+70
View File
@@ -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/<name>.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 <name>`: Archives a job via `retention-archive-jobs.py`.
- `job-unarchive <name>`: Restores an archived job to active `jobs/`.
- `bin/super-cli.py`: Exposes `box job archive <name>`, `box job unarchive <name>`, 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 <name>
box job unarchive <name>
```
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).
+157
View File
@@ -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()