#!/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" # Sidechat rotation: persistent reuse_key threads accumulate full history # and every dispatch re-sends it (cloud context), so a stale thread burns # full-thread tokens per nod. Cap counted threads by dispatch budget and # flush uncounted legacy threads past the age cap. SIDECHAT_MAX_DISPATCHES = 48 SIDECHAT_LEGACY_MAX_AGE_HOURS = 24 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 should_rotate_sidechat(record, current_title, now=None, max_dispatches=SIDECHAT_MAX_DISPATCHES, legacy_max_age_hours=SIDECHAT_LEGACY_MAX_AGE_HOURS): """Decide whether a reused sidechat must rotate to a fresh thread. Returns (rotate, reason). Rotates when the dispatch budget is spent, the rendered title moved on (daily {date} templates), or an uncounted legacy record is past the age cap. Anything unassessable (plain-UUID records, missing/unparseable age) fails open to reuse. """ now = now or datetime.now(timezone.utc) if not isinstance(record, dict): return False, "unrecorded" count = record.get("dispatch_count") if isinstance(count, int) and count >= max_dispatches: return True, f"dispatch budget spent ({count}/{max_dispatches})" stored_title = record.get("title") or "" ALLOW_SIDECHAT_TITLE_ROTATION = False if ALLOW_SIDECHAT_TITLE_ROTATION and stored_title and current_title and stored_title != current_title: return True, f"title rolled over ({stored_title} -> {current_title})" if count is None: created = record.get("created_at") if created: try: age_h = (now - datetime.fromisoformat( str(created).replace("Z", "+00:00"))).total_seconds() / 3600 except Exception: return False, "unparseable age" if age_h > legacy_max_age_hours: return True, (f"predates counting, age {age_h:.0f}h " f"over {legacy_max_age_hours}h cap") return False, "within budget" 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 # Dispatch backpressure (2026-10-09): skip jobs for frozen agents instead of # piling input-waits onto them. See tests/test_dispatch_hold.py. DISPATCH_HOLD_FILE = JOBS_DIR / "dispatch-hold.json" HOLD_WAIT_THRESHOLD = 3 HOLD_WAIT_WINDOW_MIN = 60 def dispatch_hold_reason(agent, now=None, hold_path=None, job_log_path=None): # Hold reason if dispatch to agent must be skipped, else None. # Explicit operator holds win; otherwise auto-hold after repeated waits. from datetime import timedelta now = now or datetime.now(timezone.utc) try: with open(hold_path or DISPATCH_HOLD_FILE) as f: holds = json.load(f) except (OSError, ValueError): holds = {} entry = holds.get(agent) if isinstance(holds, dict) else None if isinstance(entry, dict): until = entry.get("until") if until: try: exp = datetime.fromisoformat(until) if exp.tzinfo is None: exp = exp.replace(tzinfo=timezone.utc) except ValueError: exp = None if exp is not None and exp <= now: entry = None if entry is not None: return "explicit hold (%s)" % entry.get("reason", "operator") try: cutoff = now - timedelta(minutes=HOLD_WAIT_WINDOW_MIN) n = 0 with open(job_log_path or JOB_LOG) as f: for line in f: try: r = json.loads(line) except ValueError: continue if r.get("type") != "job_dispatch_agent_input_wait": continue if r.get("agent") != agent: continue try: ts = datetime.fromisoformat(r.get("ts", "")) except ValueError: continue if ts.tzinfo is None: ts = ts.replace(tzinfo=timezone.utc) if ts >= cutoff: n += 1 if n >= HOLD_WAIT_THRESHOLD: return "auto-hold (%d input-waits in last %dm)" % (n, HOLD_WAIT_WINDOW_MIN) except OSError: pass return None 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(): archived_file = JOBS_DIR / "archive" / f"{job_name}.json" if archived_file.exists(): print(f"Error: Job '{job_name}' is archived at {archived_file}. Unarchive before dispatch (e.g. 'box job unarchive {job_name}').", file=sys.stderr) sys.exit(1) 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) # Cryptographic attestation: only sign if explicitly requested by job config # to avoid blowing up agent context windows with massive base64 SSH signature blocks. signed_payload = None if os.environ.get("JOB_REQUIRE_SIGNATURE") == "1": dm_sign_sh = NETVM_ROOT / "bin" / "dm-sign.sh" priv_key = Path(os.path.expanduser("~/.ssh/id_ed25519")) if dm_sign_sh.exists() and priv_key.exists(): try: sign_res = subprocess.run( [str(dm_sign_sh), "--from", "super", "--key", str(priv_key), message], capture_output=True, text=True, timeout=10 ) if sign_res.returncode == 0 and "-----BEGIN SSH SIGNATURE-----" in sign_res.stdout: signed_payload = sign_res.stdout.strip() id_m = re.search(r"\[id:([a-f0-9]+)\]", signed_payload) proof_id = id_m.group(1) if id_m else None if proof_id: proof_data = { "id": proof_id, "signer": "super", "target_agent": agent, "target_conversation": target, "raw_payload": signed_payload, "ts": datetime.now(timezone.utc).isoformat() } try: import urllib.request req = urllib.request.Request( "https://crypt.muse-dev.online/proofs", data=json.dumps(proof_data).encode("utf-8"), headers={"Content-Type": "application/json", "User-Agent": "job-dispatch/1.0"}, method="POST" ) with urllib.request.urlopen(req, timeout=3) as resp: pass except Exception as pe: sys.stderr.write(f"warning: proof registration to crypt.muse-dev.online failed: {pe}\n") except Exception as se: sys.stderr.write(f"warning: dm signing failed: {se}\n") payload_to_send = signed_payload or message # Try fast hybrid gateway send if target resolves to UUID target_uuid = None try: import dm target_uuid = dm.resolve_sidechat_target(target, agent) except Exception: pass if target_uuid and re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", target_uuid.lower()): try: import muse_hybrid send_node = agent if agent in ("muse", "pip", "646", "opm") else "opm" res, err = muse_hybrid.send_message(send_node, payload_to_send, thread_id=target_uuid, wait=0) if res and not err: m_id = None id_m = re.search(r"\[id:([a-f0-9]+)\]", payload_to_send) if id_m: m_id = id_m.group(1) log_entry = { "type": "sent", "id": m_id or (res.get("reply", {}).get("message_id") if isinstance(res, dict) else "gateway"), "agent": "opm", "to": agent, "target": target, "thread_uuid": target_uuid, "transport": "gateway", "verified": True, "ts": datetime.now(timezone.utc).isoformat() } dm_log_path = NETVM_ROOT / "dm-log.jsonl" with open(dm_log_path, "a", encoding="utf-8") as lf: lf.write(json.dumps(log_entry) + "\n") return m_id or "gateway-verified" sys.stderr.write(f"warning: fast gateway send fallback: {err or res}\n") except Exception as ge: sys.stderr.write(f"warning: fast gateway send fallback: {ge}\n") if signed_payload: cmd = ([str(DM_PY), "send", "--agent", "opm", "--to", agent, "--target", target, "--raw"] + (["--allow-main-chat"] if (target == "main" and allow_main_chat) else []) + (followup_tags or []) + [signed_payload]) else: 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) # Backpressure: skip frozen agents before arming follow-ups or sending. _hold = dispatch_hold_reason(job.get("agent")) if _hold: print("Held: job %s for %s skipped (%s)." % (job_name, job.get("agent"), _hold), file=sys.stderr) log_event("job_dispatch_held", {"job_name": job_name, "agent": job.get("agent"), "reason": _hold}) sys.exit(0) # 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) # 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("sidechat", {}).get("create"): sc_cfg = job.get("sidechat", {}) sc_name = render_prompt(sc_cfg.get("name_template", "job-{job_name}-{date}"), variables) reuse_key = sc_cfg.get("reuse_key") # Check if reuse_key exists in job-sidechats.json and thread is still alive sc_state = load_sidechat_state() reused_uuid = None rotated_from = None if reuse_key and reuse_key in sc_state: val = sc_state[reuse_key] cand_uuid = val.get("thread_uuid") if isinstance(val, dict) else val if cand_uuid: rotate, reason = should_rotate_sidechat(val, sc_name) if rotate: print(f"Rotating sidechat '{reuse_key}': {reason}") log_event("job_sidechat_rotate", { "job_name": job_name, "job_id": job_id, "reuse_key": reuse_key, "old_thread": cand_uuid, "reason": reason, }) rotated_from = cand_uuid else: try: import muse_hybrid threads, err = muse_hybrid.get_threads(agent) if not err and threads: thread_ids = [t.get("session_id") for t in threads] if cand_uuid in thread_ids: reused_uuid = cand_uuid except Exception: pass if reused_uuid and isinstance(val, dict): val["dispatch_count"] = val.get("dispatch_count", 0) + 1 save_sidechat_state(sc_state) if reused_uuid: target = reused_uuid else: # Spawn a brand new sidechat/channel via fast headless gateway! channel_title = sc_name try: import muse_hybrid res, err = muse_hybrid.start_session(agent, title=channel_title) if res and not err and res.get("session_id"): new_uuid = res.get("session_id") key_to_save = reuse_key or sc_name is_persistent = bool(reuse_key) new_record = { "thread_uuid": new_uuid, "agent": agent, "title": channel_title, "type": "persistent" if is_persistent else "ephemeral", "created_at": datetime.now(timezone.utc).isoformat(), "dispatch_count": 1, } if rotated_from: new_record["rotated_from"] = rotated_from new_record["rotated_at"] = datetime.now( timezone.utc).isoformat() sc_state[key_to_save] = new_record save_sidechat_state(sc_state) target = new_uuid print(f"Spawned new sidechat channel '{channel_title}' ({new_uuid}) for {agent}") else: target = sc_name except Exception as e: sys.stderr.write(f"warning: fast gateway session-start exception ({e}), falling back to name {sc_name}\n") target = sc_name elif job.get("dm_target"): target = job.get("dm_target").strip() elif job.get("target"): target = job.get("target").strip() # Pre-dispatch approval & input check (auto-approve trusted; warn if blocked) if not dry_run: try: import approvals app_info = approvals.inspect_node_approvals(agent) if app_info.get("has_pending"): if app_info.get("is_trusted") and app_info.get("status") != "KEY_APPROVAL": print(f"Pre-dispatch: auto-approving trusted request for {agent} ({app_info.get('target')})") approvals.allow_node_approval(agent, always=True, caller="job-dispatch") else: print(f"Warning: Agent '{agent}' has untrusted pending approval ({app_info.get('target')}). Dispatch may stall.", file=sys.stderr) log_event("job_dispatch_approval_blocked", {"job_id": job_id, "agent": agent, "target": app_info.get("target")}) elif app_info.get("status") == "INPUT_WAIT": waits = app_info.get("input_waits", []) w_desc = "; ".join(w.get("task", "") for w in waits)[:80] print(f"Notice: Agent '{agent}' has task waiting for input ({w_desc}).", file=sys.stderr) log_event("job_dispatch_agent_input_wait", {"job_id": job_id, "agent": agent, "waits": w_desc}) except Exception: pass # 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, }) # Work-first envelope: executable swarm.spawn/followup.create at TOP and BOTTOM # (see bin/prompt_envelope.py). Skipped when the job sets "skip_envelope": true # (agents whose runtime lacks the enveloped tools, e.g. pip). if not job.get("skip_envelope"): import prompt_envelope rendered = prompt_envelope.wrap(job_name, job_id, agent, target, rendered) # Format as JOB DM dm_message = f"[JOB {job_id}] {rendered}" # 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()