#!/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, 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 hashlib import json import os import re 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") 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 VALID_AGENTS = ["muse", "pip", "646", "opm"] DEFAULT_PORTS = {"muse": 9410, "pip": 9420, "646": 9430, "opm": 9440} 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") def get_monitored_threads(target_agent=None): """ Build dict of threads to monitor per agent: { agent: [ {"id": "main", "name": "main"}, {"id": "", "name": ""} ] } """ agents = [target_agent] if target_agent else VALID_AGENTS threads_by_agent = {a: [{"id": "main", "name": "Main Chat"}] for a in agents} 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 if isinstance(val, dict): uuid = val.get("thread_uuid") or val.get("uuid") agent = val.get("agent", "opm") elif isinstance(val, str): uuid = val agent = "opm" else: continue if uuid and agent in threads_by_agent: # Avoid duplicate threads 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 harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run=False): """ Harvests new messages for a single thread, 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, "") # 1. Capture current URL before navigating initial_url = cdp.evaluate("window.location.href") or "https://muse.ai/" is_init_main = "/thread/" not in initial_url # 2. Navigate to target thread if not already there try: if thread_id == "main": if not is_init_main: # Dispatch Ctrl+J (modifier 2 = Control) 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)}") # Settle wait time.sleep(2.5) # 3. Scrape messages raw_messages = scrape_thread_messages(cdp, thread_id, last_wm) finally: # 4. State preservation: restore browser back to initial state 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 if not raw_messages: return [], last_wm, 0 # 5. Filter for new messages based on watermark new_messages = [] if not last_wm: # Initial run on this thread: watermark at current latest message to avoid flooding backlog new_wm = raw_messages[-1]["id"] return [], new_wm, 0 else: # Find index of last_wm 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 (older than scroll limit) # Process all visible messages that are newer than timestamp or just unread tail new_messages = raw_messages if not new_messages: return [], last_wm, 0 new_wm = new_messages[-1]["id"] job_results = 0 # 6. Ingest new messages 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() # Append to chat-history.jsonl 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 not dry_run: append_jsonl(CHAT_HISTORY_LOG, record) # Check for [RESULT ] in assistant messages if author == "assistant": m_res = re.search(r"\[RESULT\s+([A-Za-z0-9_-]+)\]\s*(.*)", text, re.S) if m_res: job_id = m_res.group(1).strip() result_text = m_res.group(2).strip() is_fail = result_text.startswith("FAILED") or result_text.startswith("UNABLE") or result_text.startswith("FAIL") 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) # Trigger pipeline chaining or next job if configured trigger_chain_next(job_id, result_text, success=not is_fail) # Check and clear pending follow-ups clear_matching_followups(followups, agent, thread_id, mid, text, dry_run) return new_messages, new_wm, job_results def trigger_chain_next(job_id, result_text, success=True): """If the completed job has on_success, on_failure, or chain_next, dispatch downstream.""" 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): """Resolve follow-up records if an assistant message is detected in the thread.""" if not followups: return modified = False for f_id, f_rec in followups.items(): if 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 if match_thread: 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} peer_ip, port = get_node_network(agent) # Wrap in CDP slot if available slot_ctx = ( cdp_slot(agent, priority=PRIORITY_LOW, timeout=10) if HAS_CDP_QUEUE else None ) try: if slot_ctx: with slot_ctx: 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 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 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) # 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()