From da59ad6edd4e67fcd610f5ab1d852acc5cb3fbbd Mon Sep 17 00:00:00 2001 From: operator Date: Mon, 5 Oct 2026 03:21:20 +0000 Subject: [PATCH] feat(harvester,jobs): fast gateway harvesting, parallel thread polling, dynamic sidechat spawning - Implement muse_hybrid fast gateway in response-harvester.py with ThreadPoolExecutor - Reduce harvest cycle time from ~1m50s to ~18s; eliminate browser tab-hopping and CDP lock contention - Add smart active thread filtering in get_monitored_threads to prune dead historical test pipes - Fix initial watermark ingestion logic so fresh sidechats process first-arrival job responses - Configure gw_wait window in dm.py (20s on reply:expected) so backend generation is not prematurely severed - Update box-http-health, box-service-health, box-deep-health to dynamic timestamped sidechats without reuse_key - Fix heartbeat job to route to dedicated heartbeat channel --- bin/dm.py | 6 +- bin/response-harvester.py | 356 ++++++++++++++++++++--------------- jobs/box-deep-health.json | 3 +- jobs/box-http-health.json | 3 +- jobs/box-service-health.json | 3 +- jobs/canary-test.json | 13 +- jobs/heartbeat.json | 2 +- 7 files changed, 214 insertions(+), 172 deletions(-) diff --git a/bin/dm.py b/bin/dm.py index 180265d..8695b58 100755 --- a/bin/dm.py +++ b/bin/dm.py @@ -578,7 +578,8 @@ def dm_send(agent, target, message, verify=True, raw=False, if is_uuid and recipient in VALID_AGENTS: try: import muse_hybrid - gw_res, gw_err = muse_hybrid.send_message(recipient, tagged, thread_id=nav_target, wait=0) + gw_wait = 20 if tags.get("reply:expected") else 10 + gw_res, gw_err = muse_hybrid.send_message(recipient, tagged, thread_id=nav_target, wait=gw_wait) if gw_res and not gw_err: thread_uuid = nav_target tags["thread"] = thread_uuid @@ -635,7 +636,8 @@ def dm_send(agent, target, message, verify=True, raw=False, os.replace(tmp_sc, SIDCHAT_MAP_FILE) # Send immediately via gateway - gw_res, gw_err = muse_hybrid.send_message(recipient, tagged, thread_id=matched_uuid, wait=0) + gw_wait = 20 if tags.get("reply:expected") else 10 + gw_res, gw_err = muse_hybrid.send_message(recipient, tagged, thread_id=matched_uuid, wait=gw_wait) if gw_res and not gw_err: thread_uuid = matched_uuid tags["thread"] = thread_uuid diff --git a/bin/response-harvester.py b/bin/response-harvester.py index c1b3079..b097d32 100755 --- a/bin/response-harvester.py +++ b/bin/response-harvester.py @@ -23,6 +23,7 @@ Usage: """ import argparse +import concurrent.futures import hashlib import json import os @@ -69,6 +70,12 @@ try: 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} @@ -161,12 +168,23 @@ def is_fail_result(result_text): def get_monitored_threads(target_agent=None): """ Build dict of threads to monitor per agent: - { agent: [ {"id": "main", "name": "main"}, {"id": "", "name": ""} ] } + { agent: [ {"id": "", "name": ""} ] } + Filters to permanent channels, threads with pending followups, or recent threads (< 3h). """ agents = [target_agent] if target_agent else VALID_AGENTS - # Sidechats-only: Do not monitor Main Chat to completely eliminate automated Main Chat DOM interaction 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(): @@ -176,20 +194,40 @@ def get_monitored_threads(target_agent=None): 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 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}) + 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 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 @@ -316,32 +354,26 @@ def scrape_thread_messages(cdp, thread_id, watermark, max_scrollbacks=3): return messages -def harvest_main_feed(cdp, agent, watermarks, followups, dry_run=False): +def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, followups, dry_run=False, feed=None): """ - 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). + 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) """ - 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 + # 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): @@ -351,6 +383,7 @@ def harvest_main_feed(cdp, agent, watermarks, followups, dry_run=False): 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: @@ -365,134 +398,6 @@ def harvest_main_feed(cdp, agent, watermarks, followups, dry_run=False): 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": - markers = list(iter_result_markers(text)) - verbs = list(iter_verb_markers(text)) - 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": "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, job_id=job_id, verb="RESULT") - for verb, job_id in verbs: - if verb == "RESULT": - continue # resolved via the result-marker path above - clear_matching_followups(followups, agent, "main", mid, text, - dry_run, job_id=job_id, verb=verb) - else: - 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, - 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, @@ -503,6 +408,8 @@ def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run "text": text, "source_ts": msg_ts, } + if feed: + record["feed"] = feed if not dry_run: append_jsonl(CHAT_HISTORY_LOG, record) @@ -526,13 +433,12 @@ def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run } 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) 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 # resolved via the result-marker path above + continue clear_matching_followups(followups, agent, thread_id, mid, text, dry_run, job_id=job_id, verb=verb) else: @@ -541,6 +447,110 @@ def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, 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 + ) + + + # Chain deduplication CHAINED_JOBS_FILE = Path(__file__).parent / "chained-jobs.json" @@ -731,6 +741,40 @@ def harvest_cycle(target_agent=None, dry_run=False, output_json=False): 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: diff --git a/jobs/box-deep-health.json b/jobs/box-deep-health.json index 5769efd..835eed4 100644 --- a/jobs/box-deep-health.json +++ b/jobs/box-deep-health.json @@ -15,8 +15,7 @@ "schedule": "0 9 * * *", "sidechat": { "create": true, - "name_template": "box-deep-health", - "reuse_key": "box-deep-health" + "name_template": "box-deep-health-{datetime}" }, "timeout": 1800 } diff --git a/jobs/box-http-health.json b/jobs/box-http-health.json index 2883730..57197d0 100644 --- a/jobs/box-http-health.json +++ b/jobs/box-http-health.json @@ -15,8 +15,7 @@ "schedule": "*/15 * * * *", "sidechat": { "create": true, - "name_template": "box-http-health", - "reuse_key": "box-http-health" + "name_template": "box-http-health-{datetime}" }, "timeout": 600 } diff --git a/jobs/box-service-health.json b/jobs/box-service-health.json index 2950f2d..14b0398 100644 --- a/jobs/box-service-health.json +++ b/jobs/box-service-health.json @@ -15,8 +15,7 @@ "schedule": "7,22,37,52 * * * *", "sidechat": { "create": true, - "name_template": "box-service-health", - "reuse_key": "box-service-health" + "name_template": "box-service-health-{datetime}" }, "timeout": 600 } \ No newline at end of file diff --git a/jobs/canary-test.json b/jobs/canary-test.json index 331b397..7e0bef2 100644 --- a/jobs/canary-test.json +++ b/jobs/canary-test.json @@ -1,14 +1,13 @@ { - "agent": "opm", - "chain_next": null, - "description": "Canary job to test the dispatcher - sends a simple DM", "name": "canary-test", - "on_failure": "alert", - "prompt_template": "This is a canary test from the job dispatcher.\nJob ID: {job_id}\nDate: {date}\n\nPlease reply to confirm you received this.", + "description": "Canary job to test sidechat spawning and main chat driving", + "agent": "646", "schedule": "manual", + "timeout": 300, + "on_failure": "alert", "sidechat": { "create": true, - "name_template": "canary-test" + "name_template": "canary-{datetime}" }, - "timeout": 300 + "prompt_template": "Canary test running in dedicated sidechat.\nJob ID: {job_id}\nTime: {datetime}\n\nPlease verify and reply with [RESULT {job_id}] OK." } \ No newline at end of file diff --git a/jobs/heartbeat.json b/jobs/heartbeat.json index eb03d9a..9f5ad3e 100644 --- a/jobs/heartbeat.json +++ b/jobs/heartbeat.json @@ -9,7 +9,7 @@ "sidechat": { "create": true, "name_template": "heartbeat", - "reuse_key": "heartbeat-opm" + "reuse_key": "heartbeat" }, "timeout": 300 } \ No newline at end of file