2026-10-04 16:23:10 +00:00
|
|
|
#!/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
|
|
|
|
|
|
2026-10-04 16:34:25 +00:00
|
|
|
try:
|
|
|
|
|
import pipeline_engine
|
|
|
|
|
HAS_PIPELINE = True
|
|
|
|
|
except ImportError:
|
|
|
|
|
HAS_PIPELINE = False
|
|
|
|
|
|
2026-10-04 16:23:10 +00:00
|
|
|
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": "<uuid>", "name": "<alias>"} ] }
|
|
|
|
|
"""
|
|
|
|
|
agents = [target_agent] if target_agent else VALID_AGENTS
|
2026-10-04 16:45:34 +00:00
|
|
|
# Sidechats-only: Do not monitor Main Chat to completely eliminate automated Main Chat DOM interaction
|
|
|
|
|
threads_by_agent = {a: [] for a in agents}
|
2026-10-04 16:23:10 +00:00
|
|
|
|
|
|
|
|
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
|
|
|
|
|
|
|
|
|
|
|
2026-10-04 16:45:34 +00:00
|
|
|
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
|
|
|
|
|
|
|
|
|
|
|
2026-10-04 16:23:10 +00:00
|
|
|
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 <job_id>] 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)
|
2026-10-04 16:34:25 +00:00
|
|
|
# Trigger pipeline chaining or next job if configured
|
|
|
|
|
trigger_chain_next(job_id, result_text, success=not is_fail)
|
2026-10-04 16:23:10 +00:00
|
|
|
|
|
|
|
|
# Check and clear pending follow-ups
|
|
|
|
|
clear_matching_followups(followups, agent, thread_id, mid, text, dry_run)
|
|
|
|
|
|
|
|
|
|
return new_messages, new_wm, job_results
|
|
|
|
|
|
|
|
|
|
|
2026-10-04 16:34:25 +00:00
|
|
|
def trigger_chain_next(job_id, result_text, success=True):
|
|
|
|
|
"""If the completed job has on_success, on_failure, or chain_next, dispatch downstream."""
|
2026-10-04 16:45:34 +00:00
|
|
|
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])
|
|
|
|
|
|
2026-10-04 16:23:10 +00:00
|
|
|
job_file = JOBS_DIR / f"{job_name}.json"
|
|
|
|
|
if not job_file.exists():
|
|
|
|
|
return
|
|
|
|
|
|
2026-10-04 16:34:25 +00:00
|
|
|
# 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
|
|
|
|
|
|
2026-10-04 16:23:10 +00:00
|
|
|
try:
|
|
|
|
|
with open(job_file, "r", encoding="utf-8") as f:
|
|
|
|
|
cfg = json.load(f)
|
2026-10-04 16:34:25 +00:00
|
|
|
|
|
|
|
|
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)
|
|
|
|
|
|
2026-10-04 16:23:10 +00:00
|
|
|
env = os.environ.copy()
|
|
|
|
|
env["CHAIN_PREV_JOB_ID"] = job_id
|
|
|
|
|
env["CHAIN_PREV_RESULT"] = result_text[:1000]
|
2026-10-04 16:34:25 +00:00
|
|
|
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)])
|
|
|
|
|
|
2026-10-04 16:23:10 +00:00
|
|
|
subprocess.Popen(
|
2026-10-04 16:34:25 +00:00
|
|
|
cmd,
|
2026-10-04 16:23:10 +00:00
|
|
|
env=env,
|
|
|
|
|
stdout=subprocess.DEVNULL,
|
|
|
|
|
stderr=subprocess.DEVNULL,
|
|
|
|
|
)
|
2026-10-04 16:34:25 +00:00
|
|
|
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")
|
2026-10-04 16:23:10 +00:00
|
|
|
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)
|
|
|
|
|
|
|
|
|
|
try:
|
2026-10-04 16:45:34 +00:00
|
|
|
slot_ctx = (
|
|
|
|
|
cdp_slot(agent, priority=PRIORITY_LOW, timeout=5)
|
|
|
|
|
if HAS_CDP_QUEUE
|
|
|
|
|
else None
|
|
|
|
|
)
|
2026-10-04 16:23:10 +00:00
|
|
|
if slot_ctx:
|
|
|
|
|
with slot_ctx:
|
|
|
|
|
cdp = CDPClient(agent, peer_ip, port, timeout=8)
|
|
|
|
|
cdp.connect()
|
|
|
|
|
try:
|
2026-10-04 16:45:34 +00:00
|
|
|
# 1. Harvest registered sidechat threads
|
2026-10-04 16:23:10 +00:00
|
|
|
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
|
2026-10-04 16:45:34 +00:00
|
|
|
|
|
|
|
|
# 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
|
2026-10-04 16:23:10 +00:00
|
|
|
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
|
2026-10-04 16:45:34 +00:00
|
|
|
|
|
|
|
|
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
|
2026-10-04 16:23:10 +00:00
|
|
|
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()
|