#!/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())