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
This commit is contained in:
@@ -578,7 +578,8 @@ def dm_send(agent, target, message, verify=True, raw=False,
|
|||||||
if is_uuid and recipient in VALID_AGENTS:
|
if is_uuid and recipient in VALID_AGENTS:
|
||||||
try:
|
try:
|
||||||
import muse_hybrid
|
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:
|
if gw_res and not gw_err:
|
||||||
thread_uuid = nav_target
|
thread_uuid = nav_target
|
||||||
tags["thread"] = thread_uuid
|
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)
|
os.replace(tmp_sc, SIDCHAT_MAP_FILE)
|
||||||
|
|
||||||
# Send immediately via gateway
|
# 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:
|
if gw_res and not gw_err:
|
||||||
thread_uuid = matched_uuid
|
thread_uuid = matched_uuid
|
||||||
tags["thread"] = thread_uuid
|
tags["thread"] = thread_uuid
|
||||||
|
|||||||
+163
-119
@@ -23,6 +23,7 @@ Usage:
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
import argparse
|
import argparse
|
||||||
|
import concurrent.futures
|
||||||
import hashlib
|
import hashlib
|
||||||
import json
|
import json
|
||||||
import os
|
import os
|
||||||
@@ -69,6 +70,12 @@ try:
|
|||||||
except ImportError:
|
except ImportError:
|
||||||
HAS_PIPELINE = False
|
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"]
|
VALID_AGENTS = ["muse", "pip", "646", "opm", "dev", "def"]
|
||||||
DEFAULT_PORTS = {"muse": 9410, "pip": 9420, "646": 9430, "opm": 9440, "def": 9450, "dev": 9460}
|
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):
|
def get_monitored_threads(target_agent=None):
|
||||||
"""
|
"""
|
||||||
Build dict of threads to monitor per agent:
|
Build dict of threads to monitor per agent:
|
||||||
{ agent: [ {"id": "main", "name": "main"}, {"id": "<uuid>", "name": "<alias>"} ] }
|
{ agent: [ {"id": "<uuid>", "name": "<alias>"} ] }
|
||||||
|
Filters to permanent channels, threads with pending followups, or recent threads (< 3h).
|
||||||
"""
|
"""
|
||||||
agents = [target_agent] if target_agent else VALID_AGENTS
|
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}
|
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]
|
state_files = [JOB_SIDECHATS_FILE, WAKE_SIDECHATS_FILE]
|
||||||
for sf in state_files:
|
for sf in state_files:
|
||||||
if not sf.exists():
|
if not sf.exists():
|
||||||
@@ -176,17 +194,37 @@ def get_monitored_threads(target_agent=None):
|
|||||||
for key, val in data.items():
|
for key, val in data.items():
|
||||||
if key.startswith("_"):
|
if key.startswith("_"):
|
||||||
continue
|
continue
|
||||||
|
created_at = None
|
||||||
if isinstance(val, dict):
|
if isinstance(val, dict):
|
||||||
uuid = val.get("thread_uuid") or val.get("uuid")
|
uuid = val.get("thread_uuid") or val.get("uuid")
|
||||||
agent = val.get("agent", "opm")
|
agent = val.get("agent", "opm")
|
||||||
|
created_at = val.get("created_at")
|
||||||
elif isinstance(val, str):
|
elif isinstance(val, str):
|
||||||
uuid = val
|
uuid = val
|
||||||
agent = "opm"
|
agent = "opm"
|
||||||
else:
|
else:
|
||||||
continue
|
continue
|
||||||
|
|
||||||
if uuid and agent in threads_by_agent:
|
if not uuid or agent not in threads_by_agent:
|
||||||
# Avoid duplicate threads
|
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]]
|
existing = [t["id"] for t in threads_by_agent[agent]]
|
||||||
if uuid not in existing:
|
if uuid not in existing:
|
||||||
threads_by_agent[agent].append({"id": uuid, "name": key})
|
threads_by_agent[agent].append({"id": uuid, "name": key})
|
||||||
@@ -316,30 +354,24 @@ def scrape_thread_messages(cdp, thread_id, watermark, max_scrollbacks=3):
|
|||||||
return messages
|
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.
|
Standard message processor:
|
||||||
Main Chat has no thread UUID and restoring previous state is expensive.
|
- Filters for new messages based on watermark
|
||||||
Therefore, this function operates opportunistically:
|
- Appends records to chat-history.jsonl
|
||||||
- If the browser is ALREADY on Main Chat (/thread/ not in URL): scrape directly with ZERO navigation.
|
- Detects [RESULT] and [VERB] markers, logging to job-log.jsonl and clearing followups
|
||||||
- If the browser is in a sidechat: do NOT navigate or disrupt the sidechat, skip silently.
|
- Returns (new_messages, new_watermark, job_results_count)
|
||||||
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:
|
if not raw_messages:
|
||||||
return [], last_wm, 0
|
return [], last_wm, 0
|
||||||
|
|
||||||
new_messages = []
|
new_messages = []
|
||||||
if not last_wm:
|
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"]
|
new_wm = raw_messages[-1]["id"]
|
||||||
return [], new_wm, 0
|
return [], new_wm, 0
|
||||||
else:
|
else:
|
||||||
@@ -351,6 +383,7 @@ def harvest_main_feed(cdp, agent, watermarks, followups, dry_run=False):
|
|||||||
if wm_idx >= 0:
|
if wm_idx >= 0:
|
||||||
new_messages = raw_messages[wm_idx + 1 :]
|
new_messages = raw_messages[wm_idx + 1 :]
|
||||||
else:
|
else:
|
||||||
|
# Watermark not found in loaded window
|
||||||
new_messages = raw_messages
|
new_messages = raw_messages
|
||||||
|
|
||||||
if not new_messages:
|
if not new_messages:
|
||||||
@@ -368,14 +401,15 @@ def harvest_main_feed(cdp, agent, watermarks, followups, dry_run=False):
|
|||||||
record = {
|
record = {
|
||||||
"ts": utcnow(),
|
"ts": utcnow(),
|
||||||
"agent": agent,
|
"agent": agent,
|
||||||
"thread_id": "main",
|
"thread_id": thread_id,
|
||||||
"thread_name": "Main Chat",
|
"thread_name": thread_name,
|
||||||
"msg_id": mid,
|
"msg_id": mid,
|
||||||
"author": author,
|
"author": author,
|
||||||
"text": text,
|
"text": text,
|
||||||
"source_ts": msg_ts,
|
"source_ts": msg_ts,
|
||||||
"feed": "main_chat_passive"
|
|
||||||
}
|
}
|
||||||
|
if feed:
|
||||||
|
record["feed"] = feed
|
||||||
if not dry_run:
|
if not dry_run:
|
||||||
append_jsonl(CHAT_HISTORY_LOG, record)
|
append_jsonl(CHAT_HISTORY_LOG, record)
|
||||||
|
|
||||||
@@ -394,57 +428,111 @@ def harvest_main_feed(cdp, agent, watermarks, followups, dry_run=False):
|
|||||||
"agent": agent,
|
"agent": agent,
|
||||||
"success": not is_fail,
|
"success": not is_fail,
|
||||||
"result_snippet": result_text[:300],
|
"result_snippet": result_text[:300],
|
||||||
"thread_id": "main",
|
"thread_id": thread_id,
|
||||||
"msg_id": mid,
|
"msg_id": mid,
|
||||||
}
|
}
|
||||||
if not dry_run:
|
if not dry_run:
|
||||||
append_jsonl(JOB_LOG, job_record)
|
append_jsonl(JOB_LOG, job_record)
|
||||||
trigger_chain_next(job_id, result_text, success=not is_fail)
|
trigger_chain_next(job_id, result_text, success=not is_fail)
|
||||||
clear_matching_followups(followups, agent, "main", mid, text,
|
clear_matching_followups(followups, agent, thread_id, mid, text,
|
||||||
dry_run, job_id=job_id, verb="RESULT")
|
dry_run, job_id=job_id, verb="RESULT")
|
||||||
for verb, job_id in verbs:
|
for verb, job_id in verbs:
|
||||||
if verb == "RESULT":
|
if verb == "RESULT":
|
||||||
continue # resolved via the result-marker path above
|
continue
|
||||||
clear_matching_followups(followups, agent, "main", mid, text,
|
clear_matching_followups(followups, agent, thread_id, mid, text,
|
||||||
dry_run, job_id=job_id, verb=verb)
|
dry_run, job_id=job_id, verb=verb)
|
||||||
else:
|
else:
|
||||||
clear_matching_followups(followups, agent, "main", mid, text, dry_run)
|
clear_matching_followups(followups, agent, thread_id, mid, text, dry_run)
|
||||||
|
|
||||||
return new_messages, new_wm, job_results
|
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):
|
def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run=False):
|
||||||
"""
|
"""
|
||||||
Harvests new messages for a single thread, preserves URL state,
|
Harvests new messages for a single thread via CDP (fallback),
|
||||||
and returns (new_messages, new_watermark, job_results_count).
|
preserves URL state, and returns (new_messages, new_watermark, job_results_count).
|
||||||
"""
|
"""
|
||||||
thread_id = thread_info["id"]
|
thread_id = thread_info["id"]
|
||||||
thread_name = thread_info["name"]
|
thread_name = thread_info["name"]
|
||||||
wm_key = f"{agent}:{thread_id}"
|
wm_key = f"{agent}:{thread_id}"
|
||||||
last_wm = watermarks.get(wm_key, "")
|
last_wm = watermarks.get(wm_key, "")
|
||||||
|
|
||||||
# 1. Capture current URL before navigating
|
|
||||||
initial_url = cdp.evaluate("window.location.href") or "https://muse.ai/"
|
initial_url = cdp.evaluate("window.location.href") or "https://muse.ai/"
|
||||||
is_init_main = "/thread/" not in initial_url
|
is_init_main = "/thread/" not in initial_url
|
||||||
|
|
||||||
# 2. Navigate to target thread if not already there
|
|
||||||
try:
|
try:
|
||||||
if thread_id == "main":
|
if thread_id == "main":
|
||||||
if not is_init_main:
|
if not is_init_main:
|
||||||
# Dispatch Ctrl+J (modifier 2 = Control)
|
|
||||||
cdp.dispatch_key("j", "KeyJ", modifiers=2)
|
cdp.dispatch_key("j", "KeyJ", modifiers=2)
|
||||||
time.sleep(2.0)
|
time.sleep(2.0)
|
||||||
else:
|
else:
|
||||||
target_url = f"https://muse.ai/thread/{thread_id}"
|
target_url = f"https://muse.ai/thread/{thread_id}"
|
||||||
if initial_url.strip() != target_url:
|
if initial_url.strip() != target_url:
|
||||||
cdp.evaluate(f"window.location.href = {json.dumps(target_url)}")
|
cdp.evaluate(f"window.location.href = {json.dumps(target_url)}")
|
||||||
# Settle wait
|
|
||||||
time.sleep(2.5)
|
time.sleep(2.5)
|
||||||
|
|
||||||
# 3. Scrape messages
|
|
||||||
raw_messages = scrape_thread_messages(cdp, thread_id, last_wm)
|
raw_messages = scrape_thread_messages(cdp, thread_id, last_wm)
|
||||||
finally:
|
finally:
|
||||||
# 4. State preservation: restore browser back to initial state
|
|
||||||
try:
|
try:
|
||||||
curr_url = cdp.evaluate("window.location.href") or ""
|
curr_url = cdp.evaluate("window.location.href") or ""
|
||||||
if is_init_main:
|
if is_init_main:
|
||||||
@@ -456,89 +544,11 @@ def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run
|
|||||||
except Exception:
|
except Exception:
|
||||||
pass
|
pass
|
||||||
|
|
||||||
if not raw_messages:
|
return process_messages(
|
||||||
return [], last_wm, 0
|
raw_messages, agent, thread_id, thread_name, last_wm, followups,
|
||||||
|
dry_run=dry_run
|
||||||
|
)
|
||||||
|
|
||||||
# 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)
|
|
||||||
|
|
||||||
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": 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)
|
|
||||||
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
|
|
||||||
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)
|
|
||||||
|
|
||||||
return new_messages, new_wm, job_results
|
|
||||||
|
|
||||||
|
|
||||||
# Chain deduplication
|
# Chain deduplication
|
||||||
@@ -731,6 +741,40 @@ def harvest_cycle(target_agent=None, dry_run=False, output_json=False):
|
|||||||
|
|
||||||
for agent, thread_list in monitored.items():
|
for agent, thread_list in monitored.items():
|
||||||
agent_stats = {"status": "ok", "threads": {}, "new_messages": 0, "job_results": 0}
|
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)
|
peer_ip, port = get_node_network(agent)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
|
|||||||
@@ -15,8 +15,7 @@
|
|||||||
"schedule": "0 9 * * *",
|
"schedule": "0 9 * * *",
|
||||||
"sidechat": {
|
"sidechat": {
|
||||||
"create": true,
|
"create": true,
|
||||||
"name_template": "box-deep-health",
|
"name_template": "box-deep-health-{datetime}"
|
||||||
"reuse_key": "box-deep-health"
|
|
||||||
},
|
},
|
||||||
"timeout": 1800
|
"timeout": 1800
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -15,8 +15,7 @@
|
|||||||
"schedule": "*/15 * * * *",
|
"schedule": "*/15 * * * *",
|
||||||
"sidechat": {
|
"sidechat": {
|
||||||
"create": true,
|
"create": true,
|
||||||
"name_template": "box-http-health",
|
"name_template": "box-http-health-{datetime}"
|
||||||
"reuse_key": "box-http-health"
|
|
||||||
},
|
},
|
||||||
"timeout": 600
|
"timeout": 600
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -15,8 +15,7 @@
|
|||||||
"schedule": "7,22,37,52 * * * *",
|
"schedule": "7,22,37,52 * * * *",
|
||||||
"sidechat": {
|
"sidechat": {
|
||||||
"create": true,
|
"create": true,
|
||||||
"name_template": "box-service-health",
|
"name_template": "box-service-health-{datetime}"
|
||||||
"reuse_key": "box-service-health"
|
|
||||||
},
|
},
|
||||||
"timeout": 600
|
"timeout": 600
|
||||||
}
|
}
|
||||||
@@ -1,14 +1,13 @@
|
|||||||
{
|
{
|
||||||
"agent": "opm",
|
|
||||||
"chain_next": null,
|
|
||||||
"description": "Canary job to test the dispatcher - sends a simple DM",
|
|
||||||
"name": "canary-test",
|
"name": "canary-test",
|
||||||
"on_failure": "alert",
|
"description": "Canary job to test sidechat spawning and main chat driving",
|
||||||
"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.",
|
"agent": "646",
|
||||||
"schedule": "manual",
|
"schedule": "manual",
|
||||||
|
"timeout": 300,
|
||||||
|
"on_failure": "alert",
|
||||||
"sidechat": {
|
"sidechat": {
|
||||||
"create": true,
|
"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."
|
||||||
}
|
}
|
||||||
+1
-1
@@ -9,7 +9,7 @@
|
|||||||
"sidechat": {
|
"sidechat": {
|
||||||
"create": true,
|
"create": true,
|
||||||
"name_template": "heartbeat",
|
"name_template": "heartbeat",
|
||||||
"reuse_key": "heartbeat-opm"
|
"reuse_key": "heartbeat"
|
||||||
},
|
},
|
||||||
"timeout": 300
|
"timeout": 300
|
||||||
}
|
}
|
||||||
Reference in New Issue
Block a user