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).
This commit is contained in:
Executable
+168
@@ -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 <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"
|
||||
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]} <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", {})
|
||||
|
||||
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()
|
||||
@@ -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}
|
||||
Reference in New Issue
Block a user