diff --git a/bin/job-dispatch.py b/bin/job-dispatch.py index b6d3e45..3398e4d 100755 --- a/bin/job-dispatch.py +++ b/bin/job-dispatch.py @@ -27,6 +27,46 @@ 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")) @@ -119,27 +159,6 @@ def create_sidechat(sender_agent, dry_run=False): print(f"Sidechat creation failed: {e}", file=sys.stderr) return False -def find_sidechat(sender_agent, name, dry_run=False): - """Find existing sidechat by name and navigate to it. Returns True if found.""" - if dry_run: - return False - # List sidechats and check if name exists - cmd = [NETVM_EXEC, sender_agent, "--", "python3", str(CHAT_API), - "--account", sender_agent, "sidechat", "list"] - try: - result = subprocess.run(cmd, capture_output=True, text=True, timeout=30) - if name.lower() not in result.stdout.lower(): - return False - # Navigate to it - nav = [NETVM_EXEC, sender_agent, "--", "python3", str(CHAT_API), - "--account", sender_agent, "sidechat", "use", name] - nav_r = subprocess.run(nav, capture_output=True, text=True, timeout=30) - if nav_r.returncode == 0: - print(f"Reusing sidechat: {name}", file=sys.stderr) - return True - return False - except Exception: - 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). @@ -205,23 +224,42 @@ def main(): sidechat_created = False if use_sidechat: + reuse_key = sidechat_cfg.get("reuse_key") name_tmpl = sidechat_cfg.get("name_template", "job-{job_name}-{date}") sc_name = render_prompt(name_tmpl, variables) - # Try reuse first (for recurring jobs), then create - if find_sidechat("opm", sc_name, dry_run=dry_run): - log_event("job_sidechat_reused", {"job_id": job_id, "sidechat_name": sc_name}) - sidechat_created = True - target = "sidechat" - else: + # 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: print(f"Creating sidechat for job {job_id}...", file=sys.stderr) sidechat_created = create_sidechat("opm", dry_run=dry_run) if sidechat_created: - log_event("job_sidechat_created", {"job_id": job_id, "sidechat_name": sc_name}) + log_event("job_sidechat_created", {"job_id": job_id, "sidechat_name": sc_name, "reuse_key": reuse_key}) print(f"Sidechat created: {sc_name}", file=sys.stderr) target = "sidechat" + # Capture UUID after send for future reuse + capture_uuid = bool(reuse_key and not dry_run) else: print(f"Warning: Failed to create sidechat, falling back to main", file=sys.stderr) target = "main" + capture_uuid = False else: target = "main" @@ -245,6 +283,20 @@ def main(): "target": "sidechat", "method": "direct", }) + # 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 # Skip the dm.py dispatch block below import sys as _sys2 _sys2.exit(0) diff --git a/bin/muse-chat-api.py b/bin/muse-chat-api.py index bf0d8d4..9cb0947 100755 --- a/bin/muse-chat-api.py +++ b/bin/muse-chat-api.py @@ -378,6 +378,11 @@ def ev1(ws, expr, await_p=False): return None +def cmd_url(ws): + """Print current browser URL.""" + url = ev(ws, "window.location.href") + print(url) + def cmd_upload(ws, filepath, message=None, dry_run=False): """Attach a file to the chat composer via CDP DOM.setFileInputFiles. @@ -478,7 +483,7 @@ def main(): p = argparse.ArgumentParser() p.add_argument('--account', required=True, choices=list(ACCOUNTS.keys()), help='Agent name (matches node, profile, ACCOUNTS.md)') - p.add_argument('command', choices=['send', 'messages', 'wait', 'approvals', 'sidechat', 'upload']) + p.add_argument('command', choices=['send', 'messages', 'wait', 'approvals', 'sidechat', 'upload', 'url']) p.add_argument('arg', nargs='*', default=[]) p.add_argument('--dry-run', action='store_true', help='upload: stage attachment without sending') @@ -524,6 +529,8 @@ def main(): else: print(f"ERROR: unknown sidechat subcommand: {sub}", file=sys.stderr) sys.exit(1) + elif args.command == 'url': + cmd_url(ws) elif args.command == 'upload': if not args.arg: print("ERROR: upload requires a filepath", file=sys.stderr) diff --git a/jobs/heartbeat.json b/jobs/heartbeat.json index 40f3e39..eb03d9a 100644 --- a/jobs/heartbeat.json +++ b/jobs/heartbeat.json @@ -1 +1,15 @@ -{"agent":"opm","chain_next":null,"description":"Heartbeat job - verifies DM system is working every 5 minutes (goes to sidechat)","name":"heartbeat","on_failure":"alert","prompt_template":"Heartbeat check from job scheduler.\nJob ID: {job_id}\nTime: {datetime}\n\nThis is an automated heartbeat. Reply with [RESULT {job_id}] OK to confirm the DM pipeline is healthy.","schedule":"*/5 * * * *","sidechat":{"create":true,"name_template":"heartbeat"},"timeout":300} \ No newline at end of file +{ + "agent": "opm", + "chain_next": null, + "description": "Heartbeat job - verifies DM system is working every 5 minutes (goes to sidechat)", + "name": "heartbeat", + "on_failure": "alert", + "prompt_template": "Heartbeat check from job scheduler.\nJob ID: {job_id}\nTime: {datetime}\n\nThis is an automated heartbeat. Reply with [RESULT {job_id}] OK to confirm the DM pipeline is healthy.", + "schedule": "*/5 * * * *", + "sidechat": { + "create": true, + "name_template": "heartbeat", + "reuse_key": "heartbeat-opm" + }, + "timeout": 300 +} \ No newline at end of file