#!/usr/bin/env python3 """ response-harvester.py — Fleet agent readback and response harvesting daemon. Monitors Chromebox agents (muse, pip, 646, opm), harvests incoming messages from Main Chat and registered sidechats, maintains persistent watermarks, appends to chat-history.jsonl, resolves pending follow-ups (verb-aware: [ACK|CLAIM] -> acknowledged, [RESULT|DECLINE|NO-ACTION] -> resolved, outcome recorded), and records [RESULT] completions in job-log.jsonl. Features: - Direct CDP over host veth interfaces (fast, no sudo needed). - cdp_queue integration with PRIORITY_LOW (never blocks operator/DMs). - URL state preservation (restores browser to initial thread/main via Ctrl+J). - Bounded scroll-back for virtualized DOM (#hatch-chat-scroll). - Per-node fault isolation (CDP errors on one node do not abort the cycle). - Dual-mode execution (--once for systemd timers/CLI, --loop for daemon). Usage: python3 response-harvester.py --once python3 response-harvester.py --loop --interval 30 python3 response-harvester.py --agent 646 --once """ import argparse import concurrent.futures import hashlib import json import os import re import ssl import subprocess import sys import time import urllib.error import urllib.request import websocket from datetime import datetime, timezone from pathlib import Path # Paths NETVM_ROOT = Path("/home/super/Projects/NetVM") BIN_DIR = NETVM_ROOT / "bin" LOGS_DIR = NETVM_ROOT / "logs" CHAT_HISTORY_LOG = LOGS_DIR / "chat-history.jsonl" WATERMARKS_FILE = NETVM_ROOT / "siphon-watermarks.json" JOB_LOG = NETVM_ROOT / "job-log.jsonl" FOLLOWUPS_FILE = NETVM_ROOT / "followups.json" JOB_SIDECHATS_FILE = NETVM_ROOT / "job-sidechats.json" WAKE_SIDECHATS_FILE = Path("/home/super/sidechat-wake/wake-sidechats.json") SWARM_FILE = NETVM_ROOT / "swarms.json" NUDGE_TRACKER_FILE = NETVM_ROOT / "conversation-nudge-tracker.json" JOBS_DIR = NETVM_ROOT / "jobs" DISPATCH_PY = BIN_DIR / "job-dispatch.py" # Ensure bin is in sys.path sys.path.insert(0, str(BIN_DIR)) try: from cdp_queue import cdp_slot, PRIORITY_LOW HAS_CDP_QUEUE = True except ImportError: HAS_CDP_QUEUE = False try: import netvm_registry HAS_REGISTRY = True except ImportError: HAS_REGISTRY = False try: import pipeline_engine HAS_PIPELINE = True except ImportError: HAS_PIPELINE = False try: import muse_hybrid HAS_MUSE_HYBRID = True except ImportError: HAS_MUSE_HYBRID = False VALID_AGENTS = ["muse", "pip", "646", "opm", "dev", "def"] DEFAULT_PORTS = {"muse": 9410, "pip": 9420, "646": 9430, "opm": 9440, "def": 9450, "dev": 9460} # Matches EVERY [RESULT ] marker in a message (use with finditer, not # search). The result text is lazy and stops before the next marker (or end of # text), so a message closing two jobs records each with its own text instead # of the first marker greedily swallowing the second. RESULT_RE = re.compile(r"\[RESULT\s+([A-Za-z0-9_-]+)\]\s*(.*?)(?=\[RESULT\s|\Z)", re.S) def iter_result_markers(text): """Yield (job_id, result_text) for every [RESULT ] marker in text.""" for m in RESULT_RE.finditer(text or ""): yield m.group(1).strip(), m.group(2).strip() # Verb markers for the digest response protocol: # [ACK|CLAIM|RESULT|DECLINE|NO-ACTION ]. # ACK/CLAIM acknowledge a digest (nudge-suppressed, NOT closed); # RESULT/DECLINE/NO-ACTION close the digest. Every verb match records # outcome= on the followup record. VERB_RE = re.compile(r"\[(ACK|CLAIM|RESULT|DECLINE|NO-ACTION)\s+([A-Za-z0-9_-]+)\]") def iter_verb_markers(text): """Yield (verb, job_id) for every [VERB ] marker in text.""" for m in VERB_RE.finditer(text or ""): yield m.group(1), m.group(2).strip() def parse_tool_calls(text): """ Extract structured tool/exec calls from assistant messages. Supports: 1. [TOOL ] or [EXEC ] 2. ```box / ```tool / ```exec JSON blocks """ calls = [] for m in re.finditer(r"\[(?:TOOL|EXEC)\s+([a-zA-Z0-9_.-]+)(?:\s+(.*?))?\]", text or ""): op = m.group(1).strip() raw_args = (m.group(2) or "").strip() args = {} if raw_args: try: args = json.loads(raw_args) except Exception: args = {"raw": raw_args} calls.append((op, args)) for m in re.finditer(r"```(?:box|tool|exec)\s*\n(.*?)```", text or "", re.DOTALL): block = m.group(1).strip() try: d = json.loads(block) if isinstance(d, dict) and "op" in d: calls.append((d["op"], d.get("args", {}))) except Exception: pass return calls def execute_agent_tool(agent, op, args, timeout=30): """ Execute tool call via local exec-constrained HTTPS daemon over Tailscale. Returns (ok, result_or_error_string). """ token = None token_file = os.path.expanduser(f"~/.exec-tokens/{agent}") if not os.path.exists(token_file): token_file = os.path.expanduser("~/.exec-server-token") if os.path.exists(token_file): try: with open(token_file) as f: token = f.read().strip() except Exception: pass url = "https://100.123.153.75:8444/exec" payload = json.dumps({"op": op, "args": args}).encode("utf-8") headers = { "Content-Type": "application/json", "User-Agent": f"Box-Harvester/{agent}", } if token: headers["Authorization"] = f"Bearer {token}" ctx = ssl.create_default_context() ctx.check_hostname = False ctx.verify_mode = ssl.CERT_NONE req = urllib.request.Request(url, data=payload, headers=headers, method="POST") try: with urllib.request.urlopen(req, context=ctx, timeout=timeout) as resp: data = resp.read().decode("utf-8") res = json.loads(data) if res.get("rc") == 0: return True, res.get("stdout", "").strip() or "OK" elif "stdout" in res: return False, (res.get("stdout") or res.get("stderr") or "failed").strip() elif "error" in res: return False, res["error"] return True, json.dumps(res) except urllib.error.HTTPError as he: try: err_body = he.read().decode("utf-8") return False, f"HTTP {he.code}: {err_body}" except Exception: return False, f"HTTP {he.code}" except Exception as e: return False, str(e) def format_tool_result_for_chat(op, raw_output): """ Format tool output into a clean, human-readable chat message instead of spewing a raw wall of unformatted JSON. """ if not raw_output: return "OK" try: data = json.loads(raw_output) except Exception: s = str(raw_output).strip() return s[:500] if len(s) > 500 else s if op == "health.check" and isinstance(data, dict) and "fleet" in data: nodes = data.get("fleet", []) all_ok = all(n.get("cdp_ok") and n.get("proc_alive") for n in nodes) up_count = sum(1 for n in nodes if n.get("cdp_ok") and n.get("proc_alive")) lines = [f"Fleet Health: {'ALL GREEN' if all_ok else 'DEGRADED'} ({up_count}/{len(nodes)} nodes online)"] for n in nodes: st = "OK" if n.get("cdp_ok") and n.get("proc_alive") else "FAIL" lat = n.get("latency_ms", 0) lines.append(f" • {n.get('node')}: {st} ({lat}ms)") return "\n".join(lines) if op == "cron.runs" and isinstance(data, dict) and "jobs" in data: jobs = data.get("jobs", []) return f"{len(jobs)} scheduled jobs configured (e.g. {', '.join(j.get('name') for j in jobs[:6])})" if op == "cron.status" and isinstance(data, dict): return f"Cron '{data.get('name')}': {'ACTIVE' if data.get('active') else 'INACTIVE'} (last result: {data.get('last_result', 'unknown')})" if op == "vars.list" and isinstance(data, dict) and "variables" in data: vars_dict = data.get("variables", {}) sample = ", ".join(f"{k}={v}" for k, v in list(vars_dict.items())[:5]) return f"{len(vars_dict)} variables: {sample}..." if op == "vars.get" and isinstance(data, dict): return f"{data.get('name')} = {data.get('value')}" if op == "files.read" and isinstance(data, dict): if not data.get("ok"): return f"Failed to read file: {data.get('error')}" path = data.get("path") lines = data.get("lines") trunc = " (first " + str(data.get("displayed_lines")) + " lines)" if data.get("truncated") else "" content = data.get("content", "").strip() return f"File `{path}` ({lines} lines total){trunc}:\n```\n{content}\n```" if op == "files.write" and isinstance(data, dict): if not data.get("ok"): return f"Failed to write file: {data.get('error')}" return f"File `{data.get('path')}` written successfully ({data.get('bytes_written')} bytes, {data.get('lines')} lines)." if op == "web.fetch" and isinstance(data, dict): if not data.get("ok"): return f"Failed to fetch {data.get('url')}: {data.get('error')}" status = data.get("status") url = data.get("url") trunc = " (truncated to 4KB)" if data.get("truncated") else "" return f"Web Fetch `{url}` (HTTP {status}){trunc}:\n```\n{data.get('body', '').strip()[:800]}\n```" if op == "service.status" and isinstance(data, dict): if not data.get("ok"): return f"Service check failed: {data.get('error')}" st = "ACTIVE" if data.get("active") else "INACTIVE" return f"Service `{data.get('service')}` is {st}.\nStatus: {data.get('status_line')}" if op == "service.restart" and isinstance(data, dict): if not data.get("ok"): return f"Service restart failed for `{data.get('service')}`: {data.get('error')}" return f"Service `{data.get('service')}` restarted successfully." # General fallback: compact JSON capped to 400 chars s = json.dumps(data) return s[:400] + "..." if len(s) > 400 else s def utcnow(): return datetime.now(timezone.utc).isoformat() def get_node_network(node): """Derive veth peer IP and CDP port from node identity.""" tag = hashlib.sha256(node.encode()).hexdigest()[:8] idx = int(tag[:3], 16) % 200 + 10 peer_ip = f"10.201.{idx}.2" port = None if HAS_REGISTRY: try: port = netvm_registry.port_for(node) except Exception: pass if not port: port = DEFAULT_PORTS.get(node, 9410) return peer_ip, port def load_json_file(path, default=None): if default is None: default = {} if not os.path.exists(path): return default try: with open(path, "r", encoding="utf-8") as f: return json.load(f) except Exception: return default def save_json_file(path, data): tmp_path = f"{path}.tmp.{os.getpid()}" with open(tmp_path, "w", encoding="utf-8") as f: json.dump(data, f, indent=2) os.replace(tmp_path, path) def append_jsonl(path, record): os.makedirs(os.path.dirname(os.path.abspath(path)), exist_ok=True) with open(path, "a", encoding="utf-8") as f: f.write(json.dumps(record) + "\n") # Case-insensitive failure/decline detection for [RESULT] text. # "declined"/"decline"/"reject" count as failures so declines do NOT # trigger onward pipeline chaining (previously only uppercase # FAILED/UNABLE/FAIL matched, so a lowercase "declined" was logged as success). FAIL_PREFIXES = ("failed", "unable", "fail", "declined", "decline", "reject", "error") def is_fail_result(result_text): t = (result_text or "").lstrip().lower() return t.startswith(FAIL_PREFIXES) def get_monitored_threads(target_agent=None): """ Build dict of threads to monitor per agent: { agent: [ {"id": "", "name": ""} ] } Filters to permanent channels, threads with pending followups, or recent threads (< 3h). """ agents = [target_agent] if target_agent else VALID_AGENTS threads_by_agent = {a: [] for a in agents} pending_threads = set() followups = load_json_file(FOLLOWUPS_FILE) for f in followups.values(): if f.get("status") in ("pending", "acknowledged"): tu = f.get("thread_uuid") if tu: pending_threads.add(tu) PERM_KEYWORDS = ("coord", "tasks", "task", "brain", "heartbeat", "sync", "audit", "main-loop") now = datetime.now(timezone.utc) state_files = [JOB_SIDECHATS_FILE, WAKE_SIDECHATS_FILE] for sf in state_files: if not sf.exists(): continue try: data = json.loads(sf.read_text(encoding="utf-8")) for key, val in data.items(): if key.startswith("_"): continue created_at = None if isinstance(val, dict): uuid = val.get("thread_uuid") or val.get("uuid") agent = val.get("agent", "opm") created_at = val.get("created_at") elif isinstance(val, str): uuid = val agent = "opm" else: continue if not uuid or agent not in threads_by_agent: continue is_perm = any(k in key.lower() for k in PERM_KEYWORDS) is_pending = uuid in pending_threads is_recent = False if created_at: try: cat = datetime.fromisoformat(created_at.replace("Z", "+00:00")) if (now - cat).total_seconds() < 10800: is_recent = True except Exception: pass else: if not (key.startswith("pipe-") or key.startswith("test-") or key.startswith("onboarding-")): is_recent = True # If already archived, exclude from active monitoring unless pending followup if isinstance(val, dict) and val.get("archived") and not is_pending: continue if not (is_perm or is_pending or is_recent): continue existing = [t["id"] for t in threads_by_agent[agent]] if uuid not in existing: threads_by_agent[agent].append({"id": uuid, "name": key}) except Exception: continue return threads_by_agent class CDPClient: """Lightweight direct CDP client over host veth.""" def __init__(self, node, peer_ip, port, timeout=10): self.node = node self.peer_ip = peer_ip self.port = port self.timeout = timeout self.ws = None self.msg_id = 0 def connect(self): url = f"http://{self.peer_ip}:{self.port}/json/list" req = urllib.request.Request(url) with urllib.request.urlopen(req, timeout=self.timeout) as resp: targets = json.load(resp) pages = [t for t in targets if t.get("type") == "page"] if not pages: raise RuntimeError(f"No page target found on CDP for {self.node}") ws_url = pages[0]["webSocketDebuggerUrl"] self.ws = websocket.create_connection(ws_url, timeout=self.timeout) def send_cmd(self, method, params=None): self.msg_id += 1 cid = self.msg_id payload = {"id": cid, "method": method, "params": params or {}} self.ws.send(json.dumps(payload)) while True: raw = self.ws.recv() data = json.loads(raw) if data.get("id") == cid: return data def evaluate(self, expr, await_promise=False): res = self.send_cmd( "Runtime.evaluate", {"expression": expr, "returnByValue": True, "awaitPromise": await_promise}, ) result = res.get("result", {}).get("result", {}) if res.get("result", {}).get("exceptionDetails"): desc = res["result"]["exceptionDetails"].get("text", "JS exception") raise RuntimeError(f"CDP eval error: {desc}") return result.get("value") def dispatch_key(self, key, code, modifiers=0): self.send_cmd( "Input.dispatchKeyEvent", { "type": "rawKeyDown", "key": key, "code": code, "modifiers": modifiers, "windowsVirtualKeyCode": 74 if code == "KeyJ" else 0, }, ) self.send_cmd( "Input.dispatchKeyEvent", { "type": "keyUp", "key": key, "code": code, "modifiers": modifiers, "windowsVirtualKeyCode": 74 if code == "KeyJ" else 0, }, ) def close(self): if self.ws: try: self.ws.close() except Exception: pass self.ws = None DOM_EXTRACT_JS = """(() => { const els = [...document.querySelectorAll('[data-message-id]')]; return els.map(m => { const id = m.getAttribute('data-message-id'); const ps = [...m.querySelectorAll('p')].map(p => (p.innerText || '').trim()).filter(Boolean); let text = ps.join('\\n'); if (!text) { text = (m.innerText || '').replace(/^(Assistant message:|User message:)\\s*/i, '').trim(); } const t = m.querySelector('time'); return { id: id, author: id.startsWith('assistant-msg') ? 'assistant' : 'user', text: text, ts: t ? (t.getAttribute('datetime') || t.innerText || null) : null }; }); })()""" def scrape_thread_messages(cdp, thread_id, watermark, max_scrollbacks=3): """Scrape messages with bounded scroll-back if watermark is out of view.""" messages = cdp.evaluate(DOM_EXTRACT_JS) or [] # If watermark exists and is already in view, or no watermark, no scroll-back needed seen_ids = {m["id"] for m in messages if m.get("id")} if watermark and watermark not in seen_ids and max_scrollbacks > 0: # Bounded scroll-back loop for _ in range(max_scrollbacks): cdp.evaluate("""(() => { const sc = document.getElementById('hatch-chat-scroll'); if (sc) sc.scrollTop = 0; })()""") time.sleep(0.8) older = cdp.evaluate(DOM_EXTRACT_JS) or [] for m in older: if m.get("id") and m["id"] not in seen_ids: messages.insert(0, m) seen_ids.add(m["id"]) if watermark in seen_ids: break return messages def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, followups, dry_run=False, feed=None): """ Standard message processor: - Filters for new messages based on watermark - Appends records to chat-history.jsonl - Detects [RESULT] and [VERB] markers, logging to job-log.jsonl and clearing followups - Returns (new_messages, new_watermark, job_results_count) """ if not raw_messages: return [], last_wm, 0 new_messages = [] if not last_wm: # If thread has <= 10 messages (e.g. newly spawned job sidechat), process them! # Only fast-forward watermark if this is an established thread with a deep backlog. if len(raw_messages) <= 10: new_messages = raw_messages else: new_wm = raw_messages[-1]["id"] return [], new_wm, 0 else: wm_idx = -1 for i, m in enumerate(raw_messages): if m["id"] == last_wm: wm_idx = i break if wm_idx >= 0: new_messages = raw_messages[wm_idx + 1 :] else: # Watermark not found in loaded window new_messages = raw_messages if not new_messages: return [], last_wm, 0 new_wm = new_messages[-1]["id"] job_results = 0 for msg in new_messages: mid = msg.get("id", "") author = msg.get("author", "unknown") text = msg.get("text", "") msg_ts = msg.get("ts") or utcnow() record = { "ts": utcnow(), "agent": agent, "thread_id": thread_id, "thread_name": thread_name, "msg_id": mid, "author": author, "text": text, "source_ts": msg_ts, } if feed: record["feed"] = feed if not dry_run: append_jsonl(CHAT_HISTORY_LOG, record) if author == "assistant": markers = list(iter_result_markers(text)) verbs = list(iter_verb_markers(text)) # Check for structured tool / command calls from the agent tool_calls = parse_tool_calls(text) if tool_calls and not dry_run: for op, args in tool_calls: if isinstance(args, dict) and "agent" not in args: args["agent"] = agent print(f"[{agent}] Executing tool '{op}' from message {mid[:8]} in thread {thread_name or thread_id[:8]}") t_ok, t_res = execute_agent_tool(agent, op, args) tool_record = { "ts": utcnow(), "type": "tool_exec", "agent": agent, "thread_id": thread_id, "op": op, "args": args, "success": t_ok, "result_snippet": str(t_res)[:300], "msg_id": mid, } append_jsonl(JOB_LOG, tool_record) # Post clean tool output back to the originating thread via fast headless gateway if thread_id and re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", thread_id.lower()): try: import muse_hybrid if t_ok: clean_msg = format_tool_result_for_chat(op, t_res) resp_text = f"Tool result (`{op}`):\n{clean_msg}" else: resp_text = f"Tool error (`{op}`): {t_res}" muse_hybrid.send_message(agent, resp_text, thread_id=thread_id, wait=0) except Exception as te: sys.stderr.write(f"warning: failed to post tool response back to thread: {te}\n") if markers or verbs: for job_id, result_text in markers: is_fail = is_fail_result(result_text) job_results += 1 job_record = { "ts": utcnow(), "type": "job_result", "job_id": job_id, "agent": agent, "success": not is_fail, "result_snippet": result_text[:300], "thread_id": thread_id, "msg_id": mid, } if not dry_run: append_jsonl(JOB_LOG, job_record) # Check if this is a swarm slot result: sw-YYYYMMDD-HHMMSS-xxxx/ if "/" in job_id and job_id.startswith("sw-"): try: s_sid, s_slot = job_id.split("/", 1) s_proc = subprocess.Popen( [sys.executable, str(BIN_DIR / "box-ctl.py"), "swarm-report", s_sid, s_slot], stdin=subprocess.PIPE, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL ) s_proc.communicate(input=json.dumps({"ok": not is_fail, "result": result_text}).encode("utf-8")) except Exception as se: sys.stderr.write(f"warning: failed to record swarm report: {se}\n") else: trigger_chain_next(job_id, result_text, success=not is_fail) if not is_fail: archive_ephemeral_thread(agent, thread_id, job_id=job_id) clear_matching_followups(followups, agent, thread_id, mid, text, dry_run, job_id=job_id, verb="RESULT") for verb, job_id in verbs: if verb == "RESULT": continue clear_matching_followups(followups, agent, thread_id, mid, text, dry_run, job_id=job_id, verb=verb) else: clear_matching_followups(followups, agent, thread_id, mid, text, dry_run) maybe_nudge_untagged_sidechat(agent, thread_id, thread_name, mid, text, dry_run=dry_run) return new_messages, new_wm, job_results def harvest_thread_gateway(agent, thread_info, watermarks, followups, dry_run=False): """ Harvests messages via fast headless gateway (muse_hybrid). Zero CDP connections, zero browser page disruption, zero URL hopping. """ if not HAS_MUSE_HYBRID: raise RuntimeError("muse_hybrid module not available") thread_id = thread_info["id"] thread_name = thread_info["name"] wm_key = f"{agent}:{thread_id}" last_wm = watermarks.get(wm_key, "") hist, err = muse_hybrid.get_history( agent, thread_id=None if thread_id == "main" else thread_id, limit=30 ) if err or hist is None: if err and ("not_found" in err or "404" in err): # Session no longer exists; skip cleanly return [], last_wm, 0 raise RuntimeError(f"muse_hybrid error for {agent}:{thread_id}: {err}") raw_messages = [] for m in hist: mid = m.get("message_id") or m.get("id") or "" role = m.get("role") or m.get("author") or "unknown" text = m.get("text", "") raw_messages.append({ "id": mid, "author": role, "text": text, "ts": m.get("ts") or utcnow(), }) return process_messages( raw_messages, agent, thread_id, thread_name, last_wm, followups, dry_run=dry_run, feed="fast_gateway" ) def harvest_main_feed(cdp, agent, watermarks, followups, dry_run=False): """ Dedicated Main Chat CDP harvester (fallback). Operates opportunistically: only scrapes if browser is ALREADY on Main Chat. """ curr_url = cdp.evaluate("window.location.href") or "" if "/thread/" in curr_url: return [], None, 0 wm_key = f"{agent}:main" last_wm = watermarks.get(wm_key, "") raw_messages = scrape_thread_messages(cdp, "main", last_wm, max_scrollbacks=0) return process_messages( raw_messages, agent, "main", "Main Chat", last_wm, followups, dry_run=dry_run, feed="main_chat_passive" ) def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run=False): """ Harvests new messages for a single thread via CDP (fallback), preserves URL state, and returns (new_messages, new_watermark, job_results_count). """ thread_id = thread_info["id"] thread_name = thread_info["name"] wm_key = f"{agent}:{thread_id}" last_wm = watermarks.get(wm_key, "") initial_url = cdp.evaluate("window.location.href") or "https://muse.ai/" is_init_main = "/thread/" not in initial_url try: if thread_id == "main": if not is_init_main: cdp.dispatch_key("j", "KeyJ", modifiers=2) time.sleep(2.0) else: target_url = f"https://muse.ai/thread/{thread_id}" if initial_url.strip() != target_url: cdp.evaluate(f"window.location.href = {json.dumps(target_url)}") time.sleep(2.5) raw_messages = scrape_thread_messages(cdp, thread_id, last_wm) finally: try: curr_url = cdp.evaluate("window.location.href") or "" if is_init_main: if "/thread/" in curr_url: cdp.dispatch_key("j", "KeyJ", modifiers=2) else: if curr_url.strip() != initial_url.strip(): cdp.evaluate(f"window.location.href = {json.dumps(initial_url)}") except Exception: pass return process_messages( raw_messages, agent, thread_id, thread_name, last_wm, followups, dry_run=dry_run ) def archive_ephemeral_thread(agent, thread_id, job_id=None): """ If thread_id belongs to an ephemeral job or one-off check, archive it via hybrid gateway and tag it as archived in job-sidechats.json. """ if not thread_id or not re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", thread_id.lower()): return # Check if thread is marked as persistent in job-sidechats.json try: if JOB_SIDECHATS_FILE.exists(): sc_state = json.loads(JOB_SIDECHATS_FILE.read_text(encoding="utf-8")) for key, val in sc_state.items(): if isinstance(val, dict) and val.get("thread_uuid") == thread_id: if val.get("type") == "persistent": return # Do not archive persistent coordinator channels val["archived"] = True val["archived_at"] = utcnow() val["archived_by_job"] = job_id JOB_SIDECHATS_FILE.write_text(json.dumps(sc_state, indent=2), encoding="utf-8") except Exception as e: sys.stderr.write(f"warning: failed to update job-sidechats state: {e}\n") # Call muse-threads.py archive via subprocess (runs in node's netns) try: helper = BIN_DIR / "muse-threads.py" subprocess.run( [sys.executable, str(helper), "archive", "--agent", agent, "--thread", thread_id], capture_output=True, text=True, timeout=15 ) print(f"[{agent}] Archived ephemeral thread {thread_id[:8]} upon completion of job {job_id or 'unknown'}") except Exception as ae: sys.stderr.write(f"warning: archive_ephemeral_thread failed: {ae}\n") SWARM_WORKER_POOL = ["dev", "def", "muse"] def reconcile_and_dispatch_swarms(dry_run=False): """ Autonomous Swarm Orchestrator: 1. Scans swarms.json for pending slots. 2. Dynamically allocates available auxiliary worker nodes (dev, def, muse). 3. Provisions ephemeral sidechat per slot and dispatches the task with [RESULT /]. 4. Upon completion of all slots, sends completion summary DM to originating coordinator. """ if not SWARM_FILE.exists() or dry_run: return try: swarms = json.loads(SWARM_FILE.read_text(encoding="utf-8")) except Exception: return modified = False now = utcnow() # Determine busy workers from running slots busy_workers = set() for sid, swarm in swarms.items(): if swarm.get("status") in ("pending", "running"): for slot in swarm.get("slots", []): if slot.get("status") == "running" and slot.get("agent_id"): busy_workers.add(slot["agent_id"]) # Process each swarm for sid, swarm in swarms.items(): st = swarm.get("status") creator = swarm.get("created_by") or "646" if creator not in ("646", "pip", "opm", "muse", "dev", "def"): creator = "646" # 1. Allocate & dispatch pending slots if st in ("pending", "running"): slots = swarm.get("slots", []) for s in slots: if s.get("status") == "pending": # Find first available worker worker = None for w in SWARM_WORKER_POOL: if w not in busy_workers: worker = w break if not worker: # Fallback round-robin across worker pool if all are busy worker = SWARM_WORKER_POOL[s["slot"] % len(SWARM_WORKER_POOL)] slot_idx = s["slot"] task_text = swarm.get("task", "execute subagent task") slot_target = f"{sid}-s{slot_idx}" # Attach worker to slot s["agent_id"] = worker s["status"] = "running" s["updated_ts"] = now busy_workers.add(worker) modified = True # Prepare task directive prompt prompt = ( f"[JOB {sid}/{slot_idx}] Task for swarm slot {slot_idx}:\n" f"{task_text}\n\n" f"Reply with [RESULT {sid}/{slot_idx}] OK or FAIL ." ) # Send to worker via dm.py (which handles sidechat creation & tracking) try: dm_cmd = [ sys.executable, str(BIN_DIR / "dm.py"), "send", "--agent", creator, "--to", worker, "--target", slot_target, prompt, ] subprocess.Popen(dm_cmd, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) print(f"[swarm] Dispatched slot {slot_idx} of {sid} to {worker} in sidechat {slot_target}") except Exception as de: sys.stderr.write(f"warning: failed to dispatch swarm slot: {de}\n") # Update swarm rollup status counts = {"pending": 0, "running": 0, "done": 0, "failed": 0, "killed": 0} for slot in slots: counts[slot.get("status", "pending")] = counts.get(slot.get("status", "pending"), 0) + 1 if counts["pending"] + counts["running"] == 0: swarm["status"] = "completed" if counts["failed"] == 0 else "partial" swarm["updated_ts"] = now modified = True elif counts["running"] or counts["done"]: swarm["status"] = "running" swarm["updated_ts"] = now modified = True # 2. Check if newly completed/partial and notify coordinator if swarm.get("status") in ("completed", "partial") and not swarm.get("notified_coordinator"): swarm["notified_coordinator"] = True modified = True done_cnt = sum(1 for sl in swarm.get("slots", []) if sl.get("status") == "done") total_cnt = len(swarm.get("slots", [])) summary_msg = ( f"[Swarm Report] Swarm {sid} ({swarm.get('label') or 'task'}) finished: " f"{done_cnt}/{total_cnt} slots successful. " f"Console: https://box.muse-dev.online/#dashboard" ) try: coord_dm = [ sys.executable, str(BIN_DIR / "dm.py"), "send", "--agent", "box", "--to", creator, "--target", "main", "--allow-main-chat", summary_msg, ] subprocess.Popen(coord_dm, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) print(f"[swarm] Delivered completion report for {sid} to coordinator {creator}") except Exception as ce: sys.stderr.write(f"warning: failed to notify swarm coordinator: {ce}\n") if modified: tmp_swarms = str(SWARM_FILE) + f".tmp.{os.getpid()}" with open(tmp_swarms, "w", encoding="utf-8") as f: json.dump(swarms, f, indent=2) os.replace(tmp_swarms, SWARM_FILE) def check_and_archive_terminal_swarms(): """Check swarms.json and auto-archive ephemeral sidechats for completed or terminal swarms.""" if not SWARM_FILE.exists() or not JOB_SIDECHATS_FILE.exists(): return try: swarms_data = json.loads(SWARM_FILE.read_text(encoding="utf-8")) sc_data = json.loads(JOB_SIDECHATS_FILE.read_text(encoding="utf-8")) except Exception: return for sid, swarm in swarms_data.items(): st = swarm.get("status") if st in ("completed", "partial", "killed"): # Check slots or matching sidechats for key, val in sc_data.items(): if isinstance(val, dict) and not val.get("archived") and val.get("type") != "persistent": if sid in key or (swarm.get("label") and swarm.get("label") in key): tu = val.get("thread_uuid") ag = val.get("agent", "opm") if tu: archive_ephemeral_thread(ag, tu, job_id=sid) def maybe_nudge_untagged_sidechat(agent, thread_id, thread_name, mid, text, dry_run=False): """ If an agent replies conversationally in a sidechat backed by a job or follow-up without providing [RESULT ] or tool directives, deliver a terse 1-turn nudge footer. """ if dry_run or not thread_id or thread_id == "main": return if not re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", thread_id.lower()): return # Look up active follow-ups or job association for this thread followups = load_json_file(FOLLOWUPS_FILE) matching_job_id = None for f_id, f_rec in followups.items(): if f_rec.get("status") in ("pending", "acknowledged") and f_rec.get("thread_uuid") == thread_id: matching_job_id = f_rec.get("job_id") or f_id break if not matching_job_id and JOB_SIDECHATS_FILE.exists(): try: sc_state = json.loads(JOB_SIDECHATS_FILE.read_text(encoding="utf-8")) for k, v in sc_state.items(): if isinstance(v, dict) and v.get("thread_uuid") == thread_id: if k.startswith(("box-", "job-", "pipe-")): matching_job_id = k break except Exception: pass if not matching_job_id: return tracker = load_json_file(NUDGE_TRACKER_FILE) rec = tracker.get(thread_id, {}) if rec.get("nudged_for_mid") == mid or rec.get("nudge_count", 0) >= 1: return # Inject terse nudge footer nudge_msg = f"[Nudge: To advance the workflow, reply with [RESULT {matching_job_id}] ]" try: import muse_hybrid print(f"[{agent}] Injecting 1-turn nudge into {thread_name or thread_id[:8]} for job {matching_job_id}") muse_hybrid.send_message(agent, nudge_msg, thread_id=thread_id, wait=0) tracker[thread_id] = { "ts": utcnow(), "job_id": matching_job_id, "nudged_for_mid": mid, "nudge_count": rec.get("nudge_count", 0) + 1, } save_json_file(NUDGE_TRACKER_FILE, tracker) except Exception as ne: sys.stderr.write(f"warning: failed to deliver conversational nudge: {ne}\n") # Chain deduplication CHAINED_JOBS_FILE = Path(__file__).parent / "chained-jobs.json" def _has_chained(job_id): try: if CHAINED_JOBS_FILE.exists(): import json as _j with open(CHAINED_JOBS_FILE) as f: return job_id in _j.load(f) except: pass return False def _mark_chained(job_id): try: import json as _j c = [] if CHAINED_JOBS_FILE.exists(): with open(CHAINED_JOBS_FILE) as f: c = _j.load(f) if job_id not in c: c.append(job_id) c = c[-1000:] with open(CHAINED_JOBS_FILE, "w") as f: _j.dump(c, f) except: pass def trigger_chain_next(job_id, result_text, success=True): """If the completed job has on_success, on_failure, or chain_next, dispatch downstream.""" if _has_chained(job_id): return _mark_chained(job_id) m = re.match(r"^(.*)-(\d{8}-\d{6}-[a-f0-9]{8})$", job_id) if m: job_name = m.group(1) else: parts = job_id.split("-") if len(parts) < 3: return job_name = "-".join(parts[:-2]) job_file = JOBS_DIR / f"{job_name}.json" if not job_file.exists(): return # Update pipeline ledger if this job belongs to an active pipeline pipeline_run_id = None step_n = 1 if HAS_PIPELINE: run_entry, step_entry = pipeline_engine.record_step_result(job_id, success, result_text) if run_entry: pipeline_run_id = run_entry.get("run_id") step_n = step_entry.get("step_n", 1) + 1 try: with open(job_file, "r", encoding="utf-8") as f: cfg = json.load(f) next_job = None if success: next_job = cfg.get("on_success") or cfg.get("chain_next") else: next_job = cfg.get("on_failure") if next_job and (JOBS_DIR / f"{next_job}.json").exists(): # Inter-step settle delay to prevent browser race conditions step_delay = int(cfg.get("step_delay", 5)) if step_delay > 0: time.sleep(step_delay) env = os.environ.copy() env["CHAIN_PREV_JOB_ID"] = job_id env["CHAIN_PREV_RESULT"] = result_text[:1000] if pipeline_run_id: env["CHAIN_PIPELINE_RUN_ID"] = pipeline_run_id env["CHAIN_STEP_N"] = str(step_n) cmd = [sys.executable, str(DISPATCH_PY), next_job] if pipeline_run_id: cmd.extend(["--pipeline-run", pipeline_run_id, "--step-n", str(step_n)]) subprocess.Popen( cmd, env=env, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, ) else: # End of chain for this pipeline run if pipeline_run_id and HAS_PIPELINE: if success: pipeline_engine.complete_pipeline(pipeline_run_id) else: pipeline_engine.fail_pipeline(pipeline_run_id, "step_failed_without_fallback") except Exception: pass def clear_matching_followups(followups, agent, thread_id, mid, text, dry_run=False, job_id=None, verb=None): """Resolve follow-up records if an assistant message is detected in the thread. Matches on thread identity (thread_uuid or target='main') OR on job_id (from a [RESULT ] reply). The job_id path works regardless of thread_uuid or target, fixing ghost followups with null thread_uuid. Also matches a main-chat reply when the sweeper recorded final_nudge_target='main' (final nudge routed to main chat). When verb is given ([ACK|CLAIM|RESULT|DECLINE|NO-ACTION ]), matching is job_id-scoped for non-RESULT verbs (a verb marker names the digest it answers, so it must not touch unrelated pending followups that merely share the thread); RESULT keeps the historical thread-or-job matching. Every verb match records outcome=. ACK/CLAIM set status 'acknowledged' (nudge-suppressed, NOT closed) instead of 'resolved', and also match already-'acknowledged' records so an ACK -> RESULT lifecycle closes correctly. """ if not followups: return modified = False for f_id, f_rec in followups.items(): # Verb replies can follow an ACK (ACK -> RESULT lifecycle), so verbs # also match 'acknowledged' records; plain replies only match pending. if verb: if f_rec.get("status") not in ("pending", "acknowledged"): continue elif f_rec.get("status") != "pending": continue if f_rec.get("recipient") != agent: continue # Match either exact thread_uuid, or target alias 'main' match_thread = False if f_rec.get("target") == "main" and thread_id == "main": match_thread = True elif f_rec.get("thread_uuid") and f_rec.get("thread_uuid") == thread_id: match_thread = True elif f_rec.get("final_nudge_target") == "main" and thread_id == "main": # C3: the sweeper routed the final nudge to main chat, so a # main-chat reply resolves even when the followup target is a # sidechat. match_thread = True # Match by job_id (from [RESULT ] or [VERB ]) -- # works regardless of thread_uuid or target. This is an ADDITIONAL # path, not a replacement. match_job = False if job_id and f_rec.get("job_id") and f_rec.get("job_id") == job_id: match_job = True # Non-RESULT verbs are job-scoped: they must not acknowledge/resolve # unrelated pending followups that merely share the thread. RESULT # keeps the historical thread-or-job matching. if verb and verb != "RESULT" and not match_job: continue if match_thread or match_job: if verb: f_rec["outcome"] = verb if verb in ("ACK", "CLAIM"): # Acknowledged: sweeper nudges stop (status != pending), but # the digest is NOT closed until a closing verb arrives. f_rec["status"] = "acknowledged" f_rec["acknowledged_at"] = utcnow() f_rec["acknowledged_by_mid"] = mid else: f_rec["status"] = "resolved" f_rec["resolved_at"] = utcnow() f_rec["resolved_by_mid"] = mid f_rec["resolved_snippet"] = text[:150] modified = True if modified and not dry_run: save_json_file(FOLLOWUPS_FILE, followups) def harvest_cycle(target_agent=None, dry_run=False, output_json=False): """Execute one full harvest cycle across agents and threads.""" watermarks = load_json_file(WATERMARKS_FILE) followups = load_json_file(FOLLOWUPS_FILE) monitored = get_monitored_threads(target_agent) cycle_stats = { "timestamp": utcnow(), "agents": {}, "total_new_messages": 0, "total_job_results": 0, } for agent, thread_list in monitored.items(): agent_stats = {"status": "ok", "threads": {}, "new_messages": 0, "job_results": 0} # 1. Fast headless gateway (zero CDP locks, zero browser navigation / tab hopping) if HAS_MUSE_HYBRID and muse_hybrid.is_node_configured(agent): try: all_targets = list(thread_list) + [{"id": "main", "name": "Main Chat"}] def _harvest_one(target_info): try: return target_info, harvest_thread_gateway( agent, target_info, watermarks, followups, dry_run ) except Exception: return target_info, ([], None, 0) with concurrent.futures.ThreadPoolExecutor(max_workers=min(len(all_targets) or 1, 5)) as executor: harvest_results = list(executor.map(_harvest_one, all_targets)) for t_info, (new_msgs, new_wm, j_res) in harvest_results: wm_key = f"{agent}:{t_info['id']}" if new_wm and not dry_run: watermarks[wm_key] = new_wm if new_msgs: agent_stats["threads"][t_info["name"]] = len(new_msgs) agent_stats["new_messages"] += len(new_msgs) agent_stats["job_results"] += j_res cycle_stats["agents"][agent] = agent_stats cycle_stats["total_new_messages"] += agent_stats["new_messages"] cycle_stats["total_job_results"] += agent_stats["job_results"] continue except Exception: pass # 2. Fallback: Direct CDP over host veth peer_ip, port = get_node_network(agent) try: slot_ctx = ( cdp_slot(agent, priority=PRIORITY_LOW, timeout=5) if HAS_CDP_QUEUE else None ) if slot_ctx: with slot_ctx: cdp = CDPClient(agent, peer_ip, port, timeout=8) cdp.connect() try: # 1. Harvest registered sidechat threads for t_info in thread_list: new_msgs, new_wm, j_res = harvest_agent_thread( cdp, agent, t_info, watermarks, followups, dry_run ) wm_key = f"{agent}:{t_info['id']}" if new_wm and not dry_run: watermarks[wm_key] = new_wm agent_stats["threads"][t_info["name"]] = len(new_msgs) agent_stats["new_messages"] += len(new_msgs) agent_stats["job_results"] += j_res # 2. Opportunistic Main Chat feed harvest (zero navigation, only if already on main) m_msgs, m_wm, m_res = harvest_main_feed( cdp, agent, watermarks, followups, dry_run ) if m_wm and not dry_run: watermarks[f"{agent}:main"] = m_wm if m_msgs: agent_stats["threads"]["Main Chat (Passive)"] = len(m_msgs) agent_stats["new_messages"] += len(m_msgs) agent_stats["job_results"] += m_res finally: cdp.close() else: cdp = CDPClient(agent, peer_ip, port, timeout=8) cdp.connect() try: for t_info in thread_list: new_msgs, new_wm, j_res = harvest_agent_thread( cdp, agent, t_info, watermarks, followups, dry_run ) wm_key = f"{agent}:{t_info['id']}" if new_wm and not dry_run: watermarks[wm_key] = new_wm agent_stats["threads"][t_info["name"]] = len(new_msgs) agent_stats["new_messages"] += len(new_msgs) agent_stats["job_results"] += j_res m_msgs, m_wm, m_res = harvest_main_feed( cdp, agent, watermarks, followups, dry_run ) if m_wm and not dry_run: watermarks[f"{agent}:main"] = m_wm if m_msgs: agent_stats["threads"]["Main Chat (Passive)"] = len(m_msgs) agent_stats["new_messages"] += len(m_msgs) agent_stats["job_results"] += m_res finally: cdp.close() except Exception as e: agent_stats["status"] = "error" agent_stats["error"] = str(e)[:200] cycle_stats["agents"][agent] = agent_stats cycle_stats["total_new_messages"] += agent_stats["new_messages"] cycle_stats["total_job_results"] += agent_stats["job_results"] if not dry_run: save_json_file(WATERMARKS_FILE, watermarks) reconcile_and_dispatch_swarms() check_and_archive_terminal_swarms() # Output formatting if output_json: print(json.dumps(cycle_stats)) else: ts_short = cycle_stats["timestamp"].split("T")[1][:8] summary_parts = [] for ag, st in cycle_stats["agents"].items(): if st["status"] == "ok": summary_parts.append(f"{ag}: {st['new_messages']} msgs ({st['job_results']} results)") else: summary_parts.append(f"{ag}: [UNREACHABLE: {st.get('error', 'err')[:40]}]") print(f"[{ts_short}Z] Harvest cycle: {', '.join(summary_parts)}") return cycle_stats def main(): parser = argparse.ArgumentParser(description="Fleet agent readback and response harvester") parser.add_argument("--once", action="store_true", help="Run once and exit (default)") parser.add_argument("--loop", action="store_true", help="Run continuously in a daemon loop") parser.add_argument("--interval", type=int, default=30, help="Interval in seconds for --loop (default 30)") parser.add_argument("--agent", choices=VALID_AGENTS, default=None, help="Harvest only specific agent") parser.add_argument("--dry-run", action="store_true", help="Scrape without persisting watermarks or logs") parser.add_argument("--json", action="store_true", help="Output summary as JSON") args = parser.parse_args() # Default to --once if --loop is not provided if not args.loop: harvest_cycle(target_agent=args.agent, dry_run=args.dry_run, output_json=args.json) return print(f"Starting response-harvester daemon (interval={args.interval}s, agent={args.agent or 'all'})...") while True: try: harvest_cycle(target_agent=args.agent, dry_run=args.dry_run, output_json=args.json) except Exception as e: print(f"ERROR in harvest loop: {e}", file=sys.stderr) time.sleep(args.interval) if __name__ == "__main__": main()