feat(relay): add subagent spawn, thread tools, and container box client
- Upgrade exec-constrained.py with subagent.spawn, thread.list, thread.view, pipeline.run ops - Grant full ops permissions to all fleet agent identities (646, pip, muse, opm) - Implement bin/box-relay.sh zero-dependency client supporting bearer and SSH signature auth - Add fast hybrid gateway path to dm.py for sub-2s verified deliveries - Fix wait=0 handling in super-cli.py subagent deployments - Add hourly check-in jobs and scheduler for 646, pip, muse - Document agent tooling and relay APIs in docs/AGENT-TOOLING.md
This commit is contained in:
Executable
+211
@@ -0,0 +1,211 @@
|
||||
#!/usr/bin/env python3
|
||||
"""job-scheduler.py — NetVM unified job scheduler.
|
||||
|
||||
Evaluates cron schedules in jobs/*.json and triggers due jobs via job-dispatch.py.
|
||||
Prevents duplicate dispatches using state watermarks in /home/super/Projects/NetVM/job-scheduler-state.json.
|
||||
|
||||
Usage:
|
||||
python3 bin/job-scheduler.py run [--dry-run]
|
||||
python3 bin/job-scheduler.py status
|
||||
"""
|
||||
|
||||
import os
|
||||
import sys
|
||||
import glob
|
||||
import json
|
||||
import fcntl
|
||||
import argparse
|
||||
import subprocess
|
||||
from datetime import datetime, timezone, timedelta
|
||||
|
||||
BASE = "/home/super/Projects/NetVM"
|
||||
JOBS_DIR = os.path.join(BASE, "jobs")
|
||||
BIN = os.path.join(BASE, "bin")
|
||||
JOB_DISPATCH = os.path.join(BIN, "job-dispatch.py")
|
||||
STATE_FILE = os.path.join(BASE, "job-scheduler-state.json")
|
||||
LOCK_FILE = os.path.join(BASE, "job-scheduler.lock")
|
||||
|
||||
|
||||
def utcnow():
|
||||
return datetime.now(timezone.utc)
|
||||
|
||||
|
||||
def parse_field(pattern, val):
|
||||
if pattern == "*":
|
||||
return True
|
||||
for part in pattern.split(","):
|
||||
if "/" in part:
|
||||
sub = part.split("/")
|
||||
step = int(sub[1])
|
||||
base = sub[0]
|
||||
start = 0 if base == "*" else int(base.split("-")[0])
|
||||
end = 59 if base == "*" else int(base.split("-")[-1])
|
||||
if start <= val <= end and (val - start) % step == 0:
|
||||
return True
|
||||
elif "-" in part:
|
||||
s, e = map(int, part.split("-"))
|
||||
if s <= val <= e:
|
||||
return True
|
||||
elif part.isdigit() and int(part) == val:
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
def cron_matches(expr, dt):
|
||||
"""Check if 5-field cron expression matches datetime dt."""
|
||||
parts = expr.strip().split()
|
||||
if len(parts) != 5:
|
||||
return False
|
||||
m, h, dom, mon, dow = parts
|
||||
dow_val = (dt.weekday() + 1) % 7 # 0=Sunday
|
||||
return (parse_field(m, dt.minute) and
|
||||
parse_field(h, dt.hour) and
|
||||
parse_field(dom, dt.day) and
|
||||
parse_field(mon, dt.month) and
|
||||
(parse_field(dow, dt.weekday() + 1) or parse_field(dow, dow_val)))
|
||||
|
||||
|
||||
def is_job_due(expr, now_dt, last_fired_dt=None):
|
||||
"""Determine if a cron job is due within the recent 5-minute sampling window."""
|
||||
if not expr or expr.strip().lower() == "manual":
|
||||
return False
|
||||
|
||||
# Check minutes in the window [now - 4min, now]
|
||||
matched_dt = None
|
||||
for offset in range(5):
|
||||
sample_dt = now_dt - timedelta(minutes=offset)
|
||||
if cron_matches(expr, sample_dt):
|
||||
matched_dt = sample_dt.replace(second=0, microsecond=0)
|
||||
break
|
||||
|
||||
if not matched_dt:
|
||||
return False
|
||||
|
||||
if last_fired_dt:
|
||||
# If fired within 4 minutes of the matched slot, skip duplicate
|
||||
diff_seconds = (now_dt - last_fired_dt).total_seconds()
|
||||
# For hourly or longer jobs, prevent re-fire within 45 minutes
|
||||
if " " in expr and expr.split()[0] != "*":
|
||||
if diff_seconds < 2700:
|
||||
return False
|
||||
elif diff_seconds < 240:
|
||||
return False
|
||||
|
||||
return True
|
||||
|
||||
|
||||
def load_state():
|
||||
try:
|
||||
with open(STATE_FILE, "r") as f:
|
||||
return json.load(f)
|
||||
except Exception:
|
||||
return {"jobs": {}, "last_run": None}
|
||||
|
||||
|
||||
def save_state(state):
|
||||
tmp = STATE_FILE + ".tmp"
|
||||
with open(tmp, "w") as f:
|
||||
json.dump(state, f, indent=2)
|
||||
os.replace(tmp, STATE_FILE)
|
||||
|
||||
|
||||
def do_run(dry_run=False):
|
||||
now = utcnow()
|
||||
now_iso = now.strftime("%Y-%m-%dT%H:%M:%SZ")
|
||||
state = load_state()
|
||||
jobs_state = state.setdefault("jobs", {})
|
||||
|
||||
job_files = sorted(glob.glob(os.path.join(JOBS_DIR, "*.json")))
|
||||
dispatched = []
|
||||
skipped = []
|
||||
|
||||
for jpath in job_files:
|
||||
try:
|
||||
with open(jpath, "r", encoding="utf-8") as f:
|
||||
data = json.load(f)
|
||||
except Exception:
|
||||
continue
|
||||
|
||||
job_name = data.get("name") or os.path.basename(jpath).replace(".json", "")
|
||||
schedule = data.get("schedule")
|
||||
if not schedule or schedule.strip().lower() == "manual":
|
||||
continue
|
||||
|
||||
j_st = jobs_state.get(job_name, {})
|
||||
last_fired_str = j_st.get("last_fired")
|
||||
last_fired_dt = None
|
||||
if last_fired_str:
|
||||
try:
|
||||
last_fired_dt = datetime.fromisoformat(last_fired_str.replace("Z", "+00:00"))
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
if is_job_due(schedule, now, last_fired_dt):
|
||||
if dry_run:
|
||||
print(f"[DRY RUN] Due job: {job_name} ({schedule})")
|
||||
dispatched.append(job_name)
|
||||
continue
|
||||
|
||||
cmd = [sys.executable, JOB_DISPATCH, job_name]
|
||||
try:
|
||||
r = subprocess.run(cmd, capture_output=True, text=True, timeout=180)
|
||||
if r.returncode == 0:
|
||||
dispatched.append(job_name)
|
||||
jobs_state[job_name] = {
|
||||
"last_fired": now_iso,
|
||||
"schedule": schedule,
|
||||
"status": "dispatched"
|
||||
}
|
||||
print(f"Dispatched job: {job_name} ({schedule})", file=sys.stderr)
|
||||
else:
|
||||
err = (r.stderr or r.stdout).strip()[-200:]
|
||||
print(f"Failed to dispatch {job_name}: {err}", file=sys.stderr)
|
||||
except Exception as e:
|
||||
print(f"Exception dispatching {job_name}: {e}", file=sys.stderr)
|
||||
else:
|
||||
skipped.append(job_name)
|
||||
|
||||
if not dry_run:
|
||||
state["last_run"] = now_iso
|
||||
save_state(state)
|
||||
|
||||
result = {
|
||||
"ok": True,
|
||||
"dispatched": dispatched,
|
||||
"dispatched_count": len(dispatched),
|
||||
"evaluated_at": now_iso
|
||||
}
|
||||
return result
|
||||
|
||||
|
||||
def main():
|
||||
p = argparse.ArgumentParser(description="NetVM Unified Job Scheduler")
|
||||
sub = p.add_subparsers(dest="cmd")
|
||||
p_run = sub.add_parser("run", help="Evaluate schedules and dispatch due jobs")
|
||||
p_run.add_argument("--dry-run", action="store_true", help="Print due jobs without dispatching")
|
||||
sub.add_parser("status", help="Show scheduler state and last run")
|
||||
|
||||
args = p.parse_args()
|
||||
cmd = args.cmd or "run"
|
||||
|
||||
if cmd == "status":
|
||||
print(json.dumps(load_state(), indent=2))
|
||||
return 0
|
||||
|
||||
if cmd == "run":
|
||||
try:
|
||||
lockfh = open(LOCK_FILE, "w")
|
||||
fcntl.flock(lockfh, fcntl.LOCK_EX | fcntl.LOCK_NB)
|
||||
except (OSError, IOError):
|
||||
print(json.dumps({"ok": False, "skipped": "already running"}))
|
||||
return 0
|
||||
|
||||
res = do_run(dry_run=getattr(args, "dry_run", False))
|
||||
print(json.dumps(res))
|
||||
return 0
|
||||
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
Reference in New Issue
Block a user