From 78dbf8f131dc6aae90fc7ebd8e3647e83da0a5ed Mon Sep 17 00:00:00 2001 From: operator-main Date: Sun, 4 Oct 2026 03:06:46 +0000 Subject: [PATCH] Add job-dispatch.py: JOB system dispatcher\n\nReads job JSON, renders prompt template, sends DM via dm.py,\nlogs to job-log.jsonl. Tested end-to-end: canary-test job\ndispatched successfully (DM 36bd80aa SENT and VERIFIED). --- bin/job-dispatch.py | 168 ++++++++++++++++++++++++++++++++++++++++++ jobs/canary-test.json | 1 + 2 files changed, 169 insertions(+) create mode 100755 bin/job-dispatch.py create mode 100644 jobs/canary-test.json diff --git a/bin/job-dispatch.py b/bin/job-dispatch.py new file mode 100755 index 0000000..91e1960 --- /dev/null +++ b/bin/job-dispatch.py @@ -0,0 +1,168 @@ +#!/usr/bin/env python3 +""" +job-dispatch.py: Dispatch a job by sending a DM to an agent. + +Usage: + job-dispatch.py [--dry-run] + +Reads /home/super/Projects/NetVM/jobs/.json, +renders the prompt template, sends DM via dm.py, logs to job-log.jsonl. + +Part of the JOB system (see docs/JOB-SPEC.md). +""" + +import sys +import os +import json +import subprocess +import uuid +from datetime import datetime, timezone +from pathlib import Path + +# Paths +NETVM_ROOT = Path("/home/super/Projects/NetVM") +JOBS_DIR = NETVM_ROOT / "jobs" +DM_PY = NETVM_ROOT / "bin" / "dm.py" +JOB_LOG = NETVM_ROOT / "job-log.jsonl" + +# Add bin to path for rate_limiter +sys.path.insert(0, str(NETVM_ROOT / "bin")) +try: + from rate_limiter import rate_limit_wait + HAS_RATE_LIMITER = True +except ImportError: + HAS_RATE_LIMITER = False + +def log_event(event_type, data): + """Append event to job-log.jsonl""" + entry = { + "ts": datetime.now(timezone.utc).isoformat(), + "type": event_type, + **data + } + with open(JOB_LOG, "a") as f: + f.write(json.dumps(entry) + "\n") + +def load_job(job_name): + """Load job JSON definition""" + # Try .json first, then .yaml (for compatibility) + job_file = JOBS_DIR / f"{job_name}.json" + if not job_file.exists(): + 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) + sys.exit(1) + print(f"Warning: YAML not supported (no PyYAML). Convert {job_file} to JSON.", file=sys.stderr) + sys.exit(1) + + with open(job_file) as f: + return json.load(f) + +def render_prompt(template, variables): + """Render prompt template with variables""" + # Simple {var} substitution + result = template + for key, value in variables.items(): + result = result.replace(f"{{{key}}}", str(value)) + return result + +def send_dm(agent, target, message, dry_run=False): + """Send DM via dm.py""" + if dry_run: + print(f"[DRY RUN] Would send to {agent} ({target}):") + print(message[:200] + "..." if len(message) > 200 else message) + return "dry-run-id" + + # Rate limit + if HAS_RATE_LIMITER: + rate_limit_wait(agent) + + cmd = [str(DM_PY), "send", "--agent", "opm", "--to", agent, + "--target", target, message] + result = subprocess.run(cmd, capture_output=True, text=True, timeout=60) + + if result.returncode != 0: + print(f"DM send failed: {result.stderr}", file=sys.stderr) + return None + + # Extract message ID from output (format: SENT [id]) + # dm.py prints the ID on success + output = result.stdout.strip() + # Try to parse ID from output + return output + +def main(): + if len(sys.argv) < 2: + print(f"Usage: {sys.argv[0]} [--dry-run]", file=sys.stderr) + sys.exit(1) + + job_name = sys.argv[1] + dry_run = "--dry-run" in sys.argv + + # Load job + job = load_job(job_name) + + # Generate job_id + job_id = f"{job_name}-{datetime.now(timezone.utc).strftime('%Y%m%d-%H%M%S')}-{uuid.uuid4().hex[:8]}" + + # Variables for template + variables = { + "job_id": job_id, + "job_name": job_name, + "date": datetime.now(timezone.utc).strftime("%Y-%m-%d"), + "datetime": datetime.now(timezone.utc).isoformat(), + } + + # Render prompt + prompt_template = job.get("prompt_template", "") + if not prompt_template: + print(f"Error: Job '{job_name}' has no prompt_template", file=sys.stderr) + sys.exit(1) + + rendered = render_prompt(prompt_template, variables) + + # Format as JOB DM + dm_message = f"[JOB {job_id}] {rendered}" + + # Get target + agent = job.get("agent", "muse") + sidechat_cfg = job.get("sidechat", {}) + + if sidechat_cfg.get("create", False): + # TODO: Create sidechat and get name + # For now, just note it in the DM + target = "main" + dm_message += f"\n\n[Sidechat: {sidechat_cfg.get('name_template', 'job-'+job_name)}]" + else: + target = "main" + + # Log job_sent + log_event("job_sent", { + "job_id": job_id, + "job_name": job_name, + "agent": agent, + "target": target, + "dry_run": dry_run, + }) + + # Send DM + msg_id = send_dm(agent, target, dm_message, dry_run=dry_run) + + if msg_id and not dry_run: + print(f"Dispatched job {job_id} to {agent} (DM: {msg_id})") + log_event("job_dispatched", { + "job_id": job_id, + "dm_id": msg_id, + }) + elif dry_run: + print(f"[DRY RUN] Job {job_id} would be dispatched to {agent}") + else: + print(f"Failed to dispatch job {job_id}", file=sys.stderr) + log_event("job_failed", { + "job_id": job_id, + "error": "dm_send_failed", + }) + sys.exit(1) + +if __name__ == "__main__": + main() diff --git a/jobs/canary-test.json b/jobs/canary-test.json new file mode 100644 index 0000000..d1705bb --- /dev/null +++ b/jobs/canary-test.json @@ -0,0 +1 @@ +{"agent":"opm","chain_next":null,"description":"Canary job to test the dispatcher - sends a simple DM","name":"canary-test","on_failure":"alert","prompt_template":"This is a canary test from the job dispatcher.\nJob ID: {job_id}\nDate: {date}\n\nPlease reply with [RESULT {job_id}] to confirm you received this.","schedule":"manual","sidechat":{"create":false},"timeout":300} \ No newline at end of file