Files

212 lines
6.4 KiB
Python
Raw Permalink Normal View History

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