Files
box/bin/job-dispatch.py
T

208 lines
6.6 KiB
Python
Raw Normal View History

#!/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
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"
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"
# 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 create_sidechat(sender_agent, dry_run=False):
"""Create a sidechat via muse-chat-api.py in sender's context.
Returns the thread URL or None on failure."""
if dry_run:
print(f"[DRY RUN] Would create sidechat for {sender_agent}")
return "https://muse.ai/thread/dry-run-uuid"
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=60)
output = result.stdout.strip() + result.stderr.strip()
# Look for "Created: <url>"
for line in output.split("\n"):
if "Created:" in line:
url = line.split("Created:")[1].strip()
return url
print(f"Sidechat create output: {output[:200]}", file=sys.stderr)
return None
except Exception as e:
print(f"Sidechat creation failed: {e}", file=sys.stderr)
return None
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", {})
sidechat_url = None
if sidechat_cfg.get("create", False):
# Create sidechat in sender's (opm's) context
print(f"Creating sidechat for job {job_id}...", file=sys.stderr)
sidechat_url = create_sidechat("opm", dry_run=dry_run)
if sidechat_url:
# Render name template for display
name_tmpl = sidechat_cfg.get("name_template", "job-{job_name}-{date}")
sc_name = render_prompt(name_tmpl, variables)
dm_message += f"\n\n[Sidechat created: {sc_name}]\nWork in this side chat. URL: {sidechat_url}"
target = "main" # DM goes to main, directs to sidechat
log_event("job_sidechat_created", {
"job_id": job_id,
"sidechat_url": sidechat_url,
"sidechat_name": sc_name,
})
else:
print(f"Warning: Failed to create sidechat, proceeding without", file=sys.stderr)
target = "main"
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()