feat(harvester): add opportunistic harvest_main_feed with zero navigation and per-node fault isolation
This commit is contained in:
+120
-5
@@ -125,7 +125,8 @@ def get_monitored_threads(target_agent=None):
|
||||
{ agent: [ {"id": "main", "name": "main"}, {"id": "<uuid>", "name": "<alias>"} ] }
|
||||
"""
|
||||
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."""
|
||||
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
|
||||
try:
|
||||
slot_ctx = (
|
||||
cdp_slot(agent, priority=PRIORITY_LOW, timeout=10)
|
||||
cdp_slot(agent, priority=PRIORITY_LOW, timeout=5)
|
||||
if HAS_CDP_QUEUE
|
||||
else None
|
||||
)
|
||||
|
||||
try:
|
||||
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()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user