#!/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). Sidechat-first policy (2026-10-04): a job that resolves to target "main" without explicit opt-in fails loudly instead of saturating main threads. Opt in via job JSON "allow_main_chat": true, or --allow-main-chat. """ import sys import re import os import json import subprocess import uuid import argparse from datetime import datetime, timezone from pathlib import Path # Paths NETVM_ROOT = Path("/home/super/Projects/NetVM") sys.path.insert(0, str(NETVM_ROOT / "bin")) try: import pipeline_engine HAS_PIPELINE = True except ImportError: HAS_PIPELINE = False JOBS_DIR = NETVM_ROOT / "jobs" DM_PY = NETVM_ROOT / "bin" / "dm.py" CHAT_API = NETVM_ROOT / "bin" / "muse-chat-api.py" NETVM_EXEC = "/home/super/Projects/NetVM/bin/netvm-exec.sh" JOB_LOG = NETVM_ROOT / "job-log.jsonl" SIDECHAT_STATE = NETVM_ROOT / "job-sidechats.json" def load_sidechat_state(): if SIDECHAT_STATE.exists(): try: return json.loads(SIDECHAT_STATE.read_text()) except Exception: return {} return {} def save_sidechat_state(state): tmp = SIDECHAT_STATE.with_suffix(".tmp") tmp.write_text(json.dumps(state, indent=2)) tmp.replace(SIDECHAT_STATE) def extract_uuid(url): m = re.search(r"/thread/([0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12})", url or "") return m.group(1) if m else None def get_current_url(agent="opm"): try: cmd = [NETVM_EXEC, agent, "--", "python3", str(CHAT_API), "--account", agent, "url"] r = subprocess.run(cmd, capture_output=True, text=True, timeout=30) if r.returncode == 0: return r.stdout.strip() except Exception: pass return None def use_sidechat_uuid(agent, thread_uuid, dry_run=False): if dry_run: return False cmd = [NETVM_EXEC, agent, "--", "python3", str(CHAT_API), "--account", agent, "sidechat", "use", thread_uuid] try: r = subprocess.run(cmd, capture_output=True, text=True, timeout=30) return r.returncode == 0 except Exception: return False # 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 # ---- follow-up tracking (DM follow-up system integration) ---------------- # Jobs opt in via a "followup" block in the job JSON: # # "followup": { # "expect_reply": true, # required: enables tracking # "timeout": "1h", # duration ("30s","15m","2h","1d") or seconds # # int; default "1h" (3600s) # "nudges": 2, # 0..10, default 2 # "escalate": "opm", # identity string, default "opm" # "route": "646-pip-coord" # optional route_id # } # # The dispatcher translates this into dm.py --tag flags using the canonical # vocabulary (box-threads/DEPLOY-DECISIONS.md). dm.py strips the tags from # delivered text and creates a dm_followup request-store record after # SENT+VERIFIED. Jobs without a followup block behave exactly as today. # # LIMITATIONS (v1): # - The heartbeat job NEVER gets follow-ups (loopback health check). # Hardcoded guard below; a followup block on heartbeat is ignored loudly. # - Sidechat sends (muse-chat-api.py direct path) do not go through dm.py, # so --tag flags cannot attach. v2 needs a record-creation path that does # not send (e.g. POST /api/box/followups, or a bl->VM queue; bl cannot # currently SSH to the VM). The dispatcher logs a warning when a # sidechat-targeted job has followup enabled. HEARTBEAT_JOB_NAME = "heartbeat" def parse_followup_duration(value): """Parse a followup timeout into seconds. Accepts int (seconds) or strings like '30s', '15m', '2h', '1d'. Returns int seconds. Raises ValueError on bad input.""" if isinstance(value, int) and not isinstance(value, bool): s = value elif isinstance(value, str): m = re.fullmatch(r"(\d+)\s*([smhd])?", value.strip().lower()) if not m: raise ValueError("bad duration %r" % (value,)) n = int(m.group(1)) unit = m.group(2) or "s" s = n * {"s": 1, "m": 60, "h": 3600, "d": 86400}[unit] else: raise ValueError("timeout must be int seconds or duration string") if not 60 <= s <= 604800: raise ValueError("timeout must be 60..604800s (1m..7d), got %d" % s) return s def build_followup_tags(followup): """Translate a job's followup block into dm.py --tag arguments. Returns a flat list like ['--tag', 'reply:timeout=3600', ...]. Returns [] if followup is falsy or expect_reply is not true. Raises ValueError on invalid config (caller logs a warning and sends the DM untagged -- the job itself must never fail over this).""" if not followup or not followup.get("expect_reply"): return [] args = [] # Bare trigger. dm.py's parse_tags splits each --tag on '='; an empty # value means "present". If the deployed dm.py requires a non-empty # value for this key, use 'reply:expected=true' instead. args += ["--tag", "reply:expected="] if "timeout" in followup: s = parse_followup_duration(followup["timeout"]) args += ["--tag", "reply:timeout=%d" % s] if "nudges" in followup: n = followup["nudges"] if not isinstance(n, int) or isinstance(n, bool) or not 0 <= n <= 10: raise ValueError("nudges must be int 0..10") args += ["--tag", "reply:nudges=%d" % n] if "escalate" in followup: e = followup["escalate"] if not isinstance(e, str) or not re.fullmatch(r"[a-z0-9_-]{1,64}", e): raise ValueError("escalate must be an identity string") args += ["--tag", "reply:escalate=%s" % e] if "route" in followup: r = followup["route"] if not isinstance(r, str) or not re.fullmatch(r"[a-z0-9_-]{1,64}", r): raise ValueError("route must be a route_id string") args += ["--tag", "route:%s" % r] # 'thread' is intentionally not settable from job JSON; it names a # specific existing thread and is filled by the dispatcher when known. return args def send_dm(agent, target, message, dry_run=False, followup_tags=None, allow_main_chat=False): """Send DM via dm.py. followup_tags: flat ['--tag', 'k=v', ...] list from build_followup_tags(), or None. allow_main_chat passes the explicit main-chat opt-in through to dm.py (sidechat-first policy).""" if dry_run: print(f"[DRY RUN] Would send to {agent} ({target}):") if followup_tags: print(f"[DRY RUN] With follow-up tags: {' '.join(followup_tags)}") 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] + (["--allow-main-chat"] if (target == "main" and allow_main_chat) else []) + (followup_tags or []) + [message]) result = subprocess.run(cmd, capture_output=True, text=True, timeout=120) if result.returncode != 0: print(f"DM send failed: {result.stderr}", file=sys.stderr) return None # Extract message ID from output (format: DM ... SENT and VERIFIED) output = result.stdout.strip() m = re.search(r"DM\s+([a-f0-9]{8})", output) return m.group(1) if m else "verified" def create_sidechat(sender_agent, dry_run=False): """Create a sidechat via muse-chat-api.py in sender's context. Returns True on success (browser now on new sidechat), False on failure.""" if dry_run: print(f"[DRY RUN] Would create sidechat for {sender_agent}") return True cmd = [NETVM_EXEC, sender_agent, "--", "python3", str(CHAT_API), "--account", sender_agent, "sidechat", "create"] try: result = subprocess.run(cmd, capture_output=True, text=True, timeout=90) output = result.stdout.strip() err = result.stderr.strip() # Success if we see "Created:" (URL may be /thread/new placeholder) if "Created:" in output: print(f"Sidechat created", file=sys.stderr) return True print(f"Sidechat create failed. stdout: {output[:300]}", file=sys.stderr) print(f"Sidechat create stderr: {err[:300]}", file=sys.stderr) print(f"Return code: {result.returncode}", file=sys.stderr) return False except Exception as e: print(f"Sidechat creation failed: {e}", file=sys.stderr) return False def send_to_current_chat(sender_agent, message, dry_run=False): """Send message to current chat via muse-chat-api.py (no navigation). Used after sidechat create - browser is already on the new chat.""" if dry_run: print(f"[DRY RUN] Would send to current chat: {message[:100]}...") return "dry-run-id" if HAS_RATE_LIMITER: rate_limit_wait(sender_agent) cmd = [NETVM_EXEC, sender_agent, "--", "python3", str(CHAT_API), "--account", sender_agent, "send", message] try: result = subprocess.run(cmd, capture_output=True, text=True, timeout=60) if result.returncode == 0: return "sent-to-sidechat" print(f"Send failed: {result.stderr[:200]}", file=sys.stderr) return None except Exception as e: print(f"Send failed: {e}", file=sys.stderr) return None def main(): p = argparse.ArgumentParser(description="Dispatch a job by sending a DM to an agent.") p.add_argument("job_name", help="Name of the job (without .json)") p.add_argument("--dry-run", action="store_true", help="Print what would be sent without sending") p.add_argument("--pipeline-run", default=os.environ.get("CHAIN_PIPELINE_RUN_ID"), help="Pipeline run ID if running as part of a pipeline") p.add_argument("--step-n", type=int, default=int(os.environ.get("CHAIN_STEP_N", "1")), help="Step sequence number in the pipeline") p.add_argument("--allow-main-chat", action="store_true", help="Explicit opt-in: allow this job to dispatch to main chat (refused by default per sidechat-first policy)") args = p.parse_args() job_name = args.job_name dry_run = args.dry_run pipeline_run_id = args.pipeline_run step_n = args.step_n # 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(), "prev_job_id": os.environ.get("CHAIN_PREV_JOB_ID", ""), "prev_result": os.environ.get("CHAIN_PREV_RESULT", ""), "pipeline_run_id": pipeline_run_id or "", "step_n": step_n, } # Follow-up tracking (opt-in via job JSON "followup" block; see helpers). # The heartbeat job is a loopback health check and must never be tracked. followup_cfg = job.get("followup") followup_tags = [] if followup_cfg: if job_name == HEARTBEAT_JOB_NAME: print(f"Warning: job '{job_name}' must not use follow-up " f"tracking (loopback); ignoring followup block", file=sys.stderr) log_event("job_followup_skipped", {"job_id": job_id, "reason": "heartbeat_loopback"}) else: try: followup_tags = build_followup_tags(followup_cfg) if followup_tags: log_event("job_followup_armed", {"job_id": job_id, "tags": followup_tags}) except ValueError as e: print(f"Warning: invalid followup block: {e}; " f"sending untagged", file=sys.stderr) log_event("job_followup_invalid", {"job_id": job_id, "error": str(e)}) # 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) # Standard completion envelope: guarantee agents know how to report completion if "[RESULT" not in rendered: rendered = rendered.rstrip() + f"\n\nWhen finished, end your response with:\n[RESULT {job_id}] : " # Format as JOB DM dm_message = f"[JOB {job_id}] {rendered}" # Determine recipient agent and target agent = job.get("agent", "muse") target = "main" if pipeline_run_id: if HAS_PIPELINE: p_entry = pipeline_engine.get_pipeline(pipeline_run_id) if not p_entry: p_entry = pipeline_engine.create_pipeline(job_name, run_id=pipeline_run_id) target = p_entry.get("target") or f"pipe-{pipeline_run_id.split('-')[-1]}" else: target = f"pipe-{pipeline_run_id.split('-')[-1]}" elif job.get("dm_target"): target = job.get("dm_target").strip() elif job.get("target"): target = job.get("target").strip() elif job.get("sidechat", {}).get("create"): sc_cfg = job.get("sidechat", {}) sc_name = render_prompt(sc_cfg.get("name_template", "job-{job_name}-{date}"), variables) target = sc_name # Sidechat-first policy (2026-10-04): refuse to dispatch to main chat # unless the job explicitly opts in. Never fall back to main silently. allow_main = bool(job.get("allow_main_chat")) or args.allow_main_chat if target == "main" and not allow_main: print(f"ERROR: Job '{job_name}' resolves to main chat; refusing by sidechat-first policy. " f"Set a sidechat target (dm_target/target/sidechat.create) in the job JSON, " f"set \"allow_main_chat\": true, or pass --allow-main-chat.", file=sys.stderr) log_event("job_failed", { "job_id": job_id, "error": "main_chat_blocked_by_policy", "pipeline_run_id": pipeline_run_id, }) sys.exit(1) # Log job_sent log_event("job_sent", { "job_id": job_id, "job_name": job_name, "agent": agent, "target": target, "dry_run": dry_run, "pipeline_run_id": pipeline_run_id, "step_n": step_n, }) # Dispatch via dm.py (handles main or sidechat with auto-provisioning and verification) msg_id = send_dm(agent, target, dm_message, dry_run=dry_run, followup_tags=followup_tags, allow_main_chat=allow_main) if msg_id and not dry_run: print(f"Dispatched job {job_id} to {agent}/{target} (DM: {msg_id})") log_event("job_dispatched", { "job_id": job_id, "dm_id": msg_id, "pipeline_run_id": pipeline_run_id, "step_n": step_n, }) if pipeline_run_id and HAS_PIPELINE: pipeline_engine.record_step_dispatch( pipeline_run_id, step_n, job_name, job_id, agent, target, dm_id=msg_id ) sc_state = load_sidechat_state() if target in sc_state: val = sc_state[target] t_uuid = val.get("thread_uuid") if isinstance(val, dict) else val if t_uuid: pipeline_engine.update_pipeline_thread(pipeline_run_id, t_uuid) elif dry_run: print(f"[DRY RUN] Job {job_id} would be dispatched to {agent}/{target}") else: print(f"Failed to dispatch job {job_id}", file=sys.stderr) log_event("job_failed", { "job_id": job_id, "error": "dm_send_failed", "pipeline_run_id": pipeline_run_id, }) if pipeline_run_id and HAS_PIPELINE: pipeline_engine.fail_pipeline(pipeline_run_id, "dm_send_failed") sys.exit(1) if __name__ == "__main__": main()