UUID-based sidechat reuse for heartbeat\n\n- Add url command to muse-chat-api.py\n- Dispatcher: reuse_key -> thread UUID mapping in job-sidechats.json\n- Capture UUID after first send, reuse on subsequent runs\n- Fixes multi-spawn bug (was matching by auto-generated title)
This commit is contained in:
+80
-28
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
+15
-1
@@ -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}
|
||||
{
|
||||
"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
|
||||
}
|
||||
Reference in New Issue
Block a user