d6ca600394
- 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
212 lines
6.4 KiB
Python
Executable File
212 lines
6.4 KiB
Python
Executable File
#!/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())
|