2026-10-04 03:06:46 +00:00
|
|
|
#!/usr/bin/env python3
|
|
|
|
|
"""
|
|
|
|
|
job-dispatch.py: Dispatch a job by sending a DM to an agent.
|
|
|
|
|
|
|
|
|
|
Usage:
|
|
|
|
|
job-dispatch.py <job_name> [--dry-run]
|
|
|
|
|
|
|
|
|
|
Reads /home/super/Projects/NetVM/jobs/<job_name>.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
|
2026-10-04 03:11:11 +00:00
|
|
|
import re
|
2026-10-04 03:06:46 +00:00
|
|
|
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"
|
2026-10-04 03:09:42 +00:00
|
|
|
CHAT_API = NETVM_ROOT / "bin" / "muse-chat-api.py"
|
|
|
|
|
NETVM_EXEC = "/home/super/Projects/NetVM/bin/netvm-exec.sh"
|
2026-10-04 03:06:46 +00:00
|
|
|
JOB_LOG = NETVM_ROOT / "job-log.jsonl"
|
2026-10-04 04:03:28 +00:00
|
|
|
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
|
2026-10-04 03:06:46 +00:00
|
|
|
|
|
|
|
|
# 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
|
|
|
|
|
|
2026-10-04 03:09:42 +00:00
|
|
|
def create_sidechat(sender_agent, dry_run=False):
|
|
|
|
|
"""Create a sidechat via muse-chat-api.py in sender's context.
|
2026-10-04 03:23:34 +00:00
|
|
|
Returns True on success (browser now on new sidechat), False on failure."""
|
2026-10-04 03:09:42 +00:00
|
|
|
if dry_run:
|
|
|
|
|
print(f"[DRY RUN] Would create sidechat for {sender_agent}")
|
2026-10-04 03:23:34 +00:00
|
|
|
return True
|
2026-10-04 03:09:42 +00:00
|
|
|
|
|
|
|
|
cmd = [NETVM_EXEC, sender_agent, "--", "python3", str(CHAT_API),
|
|
|
|
|
"--account", sender_agent, "sidechat", "create"]
|
|
|
|
|
try:
|
2026-10-04 03:23:34 +00:00
|
|
|
result = subprocess.run(cmd, capture_output=True, text=True, timeout=90)
|
|
|
|
|
output = result.stdout.strip()
|
2026-10-04 03:29:21 +00:00
|
|
|
err = result.stderr.strip()
|
2026-10-04 03:23:34 +00:00
|
|
|
# Success if we see "Created:" (URL may be /thread/new placeholder)
|
|
|
|
|
if "Created:" in output:
|
|
|
|
|
print(f"Sidechat created", file=sys.stderr)
|
|
|
|
|
return True
|
2026-10-04 03:29:21 +00:00
|
|
|
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)
|
2026-10-04 03:23:34 +00:00
|
|
|
return False
|
2026-10-04 03:09:42 +00:00
|
|
|
except Exception as e:
|
|
|
|
|
print(f"Sidechat creation failed: {e}", file=sys.stderr)
|
2026-10-04 03:23:34 +00:00
|
|
|
return False
|
|
|
|
|
|
2026-10-04 03:40:25 +00:00
|
|
|
|
2026-10-04 03:23:34 +00:00
|
|
|
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)
|
2026-10-04 03:09:42 +00:00
|
|
|
return None
|
|
|
|
|
|
2026-10-04 03:06:46 +00:00
|
|
|
def main():
|
|
|
|
|
if len(sys.argv) < 2:
|
|
|
|
|
print(f"Usage: {sys.argv[0]} <job_name> [--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", {})
|
2026-10-04 03:09:42 +00:00
|
|
|
sidechat_url = None
|
2026-10-04 03:06:46 +00:00
|
|
|
|
2026-10-04 03:23:34 +00:00
|
|
|
use_sidechat = sidechat_cfg.get("create", False)
|
|
|
|
|
sidechat_created = False
|
|
|
|
|
|
|
|
|
|
if use_sidechat:
|
2026-10-04 04:03:28 +00:00
|
|
|
reuse_key = sidechat_cfg.get("reuse_key")
|
2026-10-04 03:40:25 +00:00
|
|
|
name_tmpl = sidechat_cfg.get("name_template", "job-{job_name}-{date}")
|
|
|
|
|
sc_name = render_prompt(name_tmpl, variables)
|
2026-10-04 04:03:28 +00:00
|
|
|
# UUID-based reuse: check state file for existing thread
|
|
|
|
|
state = load_sidechat_state() if reuse_key else {}
|
|
|
|
|
stored_uuid = state.get(reuse_key) if reuse_key else None
|
|
|
|
|
reused = False
|
|
|
|
|
if stored_uuid and not dry_run:
|
|
|
|
|
print(f"Trying reuse of thread {stored_uuid}...", file=sys.stderr)
|
|
|
|
|
if use_sidechat_uuid("opm", stored_uuid, dry_run=dry_run):
|
|
|
|
|
# Verify we're actually on the right thread
|
|
|
|
|
cur_url = get_current_url("opm")
|
|
|
|
|
if stored_uuid in (cur_url or ""):
|
|
|
|
|
log_event("job_sidechat_reused", {"job_id": job_id, "thread_uuid": stored_uuid, "reuse_key": reuse_key})
|
|
|
|
|
print(f"Reusing sidechat thread: {stored_uuid}", file=sys.stderr)
|
|
|
|
|
sidechat_created = True
|
|
|
|
|
target = "sidechat"
|
|
|
|
|
reused = True
|
|
|
|
|
capture_uuid = False
|
|
|
|
|
else:
|
|
|
|
|
print(f"UUID mismatch after use, creating new", file=sys.stderr)
|
|
|
|
|
else:
|
|
|
|
|
print(f"Stored thread unreachable, creating new", file=sys.stderr)
|
|
|
|
|
if not reused:
|
2026-10-04 03:40:25 +00:00
|
|
|
print(f"Creating sidechat for job {job_id}...", file=sys.stderr)
|
|
|
|
|
sidechat_created = create_sidechat("opm", dry_run=dry_run)
|
|
|
|
|
if sidechat_created:
|
2026-10-04 04:03:28 +00:00
|
|
|
log_event("job_sidechat_created", {"job_id": job_id, "sidechat_name": sc_name, "reuse_key": reuse_key})
|
2026-10-04 03:40:25 +00:00
|
|
|
print(f"Sidechat created: {sc_name}", file=sys.stderr)
|
|
|
|
|
target = "sidechat"
|
2026-10-04 04:03:28 +00:00
|
|
|
# Capture UUID after send for future reuse
|
|
|
|
|
capture_uuid = bool(reuse_key and not dry_run)
|
2026-10-04 03:40:25 +00:00
|
|
|
else:
|
|
|
|
|
print(f"Warning: Failed to create sidechat, falling back to main", file=sys.stderr)
|
|
|
|
|
target = "main"
|
2026-10-04 04:03:28 +00:00
|
|
|
capture_uuid = False
|
2026-10-04 03:06:46 +00:00
|
|
|
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,
|
|
|
|
|
})
|
|
|
|
|
|
2026-10-04 03:23:34 +00:00
|
|
|
# Send DM: to sidechat via direct API, or to main via dm.py
|
|
|
|
|
if use_sidechat and sidechat_created:
|
|
|
|
|
msg_id = send_to_current_chat("opm", dm_message, dry_run=dry_run)
|
|
|
|
|
# Log as dispatched (no dm.py ID, but sent)
|
|
|
|
|
if msg_id and not dry_run:
|
|
|
|
|
print(f"Dispatched job {job_id} to sidechat (direct send)")
|
|
|
|
|
log_event("job_dispatched", {
|
|
|
|
|
"job_id": job_id,
|
|
|
|
|
"target": "sidechat",
|
|
|
|
|
"method": "direct",
|
|
|
|
|
})
|
2026-10-04 04:03:28 +00:00
|
|
|
# Capture thread UUID for reuse_key mapping
|
|
|
|
|
if capture_uuid and reuse_key:
|
|
|
|
|
import time as _t
|
|
|
|
|
for _ in range(15):
|
|
|
|
|
_t.sleep(1)
|
|
|
|
|
cur = get_current_url("opm")
|
|
|
|
|
thread_uuid = extract_uuid(cur)
|
|
|
|
|
if thread_uuid:
|
|
|
|
|
state = load_sidechat_state()
|
|
|
|
|
state[reuse_key] = thread_uuid
|
|
|
|
|
save_sidechat_state(state)
|
|
|
|
|
log_event("job_sidechat_mapped", {"reuse_key": reuse_key, "thread_uuid": thread_uuid})
|
|
|
|
|
print(f"Mapped reuse_key {reuse_key} -> {thread_uuid}", file=sys.stderr)
|
|
|
|
|
break
|
2026-10-04 03:23:34 +00:00
|
|
|
# Skip the dm.py dispatch block below
|
|
|
|
|
import sys as _sys2
|
|
|
|
|
_sys2.exit(0)
|
|
|
|
|
else:
|
|
|
|
|
msg_id = send_dm(agent, target, dm_message, dry_run=dry_run)
|
2026-10-04 03:06:46 +00:00
|
|
|
|
|
|
|
|
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()
|