From 2fa2955fb3bbd64235751bc2fa5683737301e893 Mon Sep 17 00:00:00 2001 From: operator Date: Sun, 4 Oct 2026 16:45:34 +0000 Subject: [PATCH] feat(harvester): add opportunistic harvest_main_feed with zero navigation and per-node fault isolation --- bin/response-harvester.py | 139 ++++++++++++++++++++++++++++++++++---- 1 file changed, 127 insertions(+), 12 deletions(-) diff --git a/bin/response-harvester.py b/bin/response-harvester.py index 2b72422..d78f726 100755 --- a/bin/response-harvester.py +++ b/bin/response-harvester.py @@ -125,7 +125,8 @@ def get_monitored_threads(target_agent=None): { 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} + # Sidechats-only: Do not monitor Main Chat to completely eliminate automated Main Chat DOM interaction + threads_by_agent = {a: [] for a in agents} state_files = [JOB_SIDECHATS_FILE, WAKE_SIDECHATS_FILE] for sf in state_files: @@ -276,6 +277,95 @@ def scrape_thread_messages(cdp, thread_id, watermark, max_scrollbacks=3): return messages +def harvest_main_feed(cdp, agent, watermarks, followups, dry_run=False): + """ + Dedicated Main Chat harvester function. + Main Chat has no thread UUID and restoring previous state is expensive. + Therefore, this function operates opportunistically: + - If the browser is ALREADY on Main Chat (/thread/ not in URL): scrape directly with ZERO navigation. + - If the browser is in a sidechat: do NOT navigate or disrupt the sidechat, skip silently. + Returns (new_messages, new_watermark, job_results_count). + """ + curr_url = cdp.evaluate("window.location.href") or "" + if "/thread/" in curr_url: + # Browser is parked in a sidechat; do not disrupt the agent's work or mutate URL + return [], None, 0 + + wm_key = f"{agent}:main" + last_wm = watermarks.get(wm_key, "") + + # Scrape only current view (max_scrollbacks=0) to prevent virtual DOM layout inflation + raw_messages = scrape_thread_messages(cdp, "main", last_wm, max_scrollbacks=0) + if not raw_messages: + return [], last_wm, 0 + + new_messages = [] + if not last_wm: + 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: + 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": "main", + "thread_name": "Main Chat", + "msg_id": mid, + "author": author, + "text": text, + "source_ts": msg_ts, + "feed": "main_chat_passive" + } + if not dry_run: + append_jsonl(CHAT_HISTORY_LOG, record) + + 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": "main", + "msg_id": mid, + } + if not dry_run: + append_jsonl(JOB_LOG, job_record) + trigger_chain_next(job_id, result_text, success=not is_fail) + + clear_matching_followups(followups, agent, "main", mid, text, dry_run) + + return new_messages, new_wm, job_results + + def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run=False): """ Harvests new messages for a single thread, preserves URL state, @@ -401,10 +491,15 @@ def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run 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]) + 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 @@ -509,19 +604,18 @@ def harvest_cycle(target_agent=None, dry_run=False, output_json=False): 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: + 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 @@ -532,6 +626,17 @@ def harvest_cycle(target_agent=None, dry_run=False, output_json=False): 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: @@ -548,6 +653,16 @@ def harvest_cycle(target_agent=None, dry_run=False, output_json=False): 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()