Files
box/bin/response-harvester.py
T

851 lines
31 KiB
Python
Executable File

#!/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 (verb-aware: [ACK|CLAIM] ->
acknowledged, [RESULT|DECLINE|NO-ACTION] -> resolved, outcome recorded), 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
try:
import pipeline_engine
HAS_PIPELINE = True
except ImportError:
HAS_PIPELINE = False
VALID_AGENTS = ["muse", "pip", "646", "opm", "dev", "def"]
DEFAULT_PORTS = {"muse": 9410, "pip": 9420, "646": 9430, "opm": 9440, "def": 9450, "dev": 9460}
# Matches EVERY [RESULT <job_id>] marker in a message (use with finditer, not
# search). The result text is lazy and stops before the next marker (or end of
# text), so a message closing two jobs records each with its own text instead
# of the first marker greedily swallowing the second.
RESULT_RE = re.compile(r"\[RESULT\s+([A-Za-z0-9_-]+)\]\s*(.*?)(?=\[RESULT\s|\Z)", re.S)
def iter_result_markers(text):
"""Yield (job_id, result_text) for every [RESULT <job_id>] marker in text."""
for m in RESULT_RE.finditer(text or ""):
yield m.group(1).strip(), m.group(2).strip()
# Verb markers for the digest response protocol:
# [ACK|CLAIM|RESULT|DECLINE|NO-ACTION <job_id>].
# ACK/CLAIM acknowledge a digest (nudge-suppressed, NOT closed);
# RESULT/DECLINE/NO-ACTION close the digest. Every verb match records
# outcome=<verb> on the followup record.
VERB_RE = re.compile(r"\[(ACK|CLAIM|RESULT|DECLINE|NO-ACTION)\s+([A-Za-z0-9_-]+)\]")
def iter_verb_markers(text):
"""Yield (verb, job_id) for every [VERB <job_id>] marker in text."""
for m in VERB_RE.finditer(text or ""):
yield m.group(1), m.group(2).strip()
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")
# Case-insensitive failure/decline detection for [RESULT] text.
# "declined"/"decline"/"reject" count as failures so declines do NOT
# trigger onward pipeline chaining (previously only uppercase
# FAILED/UNABLE/FAIL matched, so a lowercase "declined" was logged as success).
FAIL_PREFIXES = ("failed", "unable", "fail", "declined", "decline", "reject", "error")
def is_fail_result(result_text):
t = (result_text or "").lstrip().lower()
return t.startswith(FAIL_PREFIXES)
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
# 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:
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
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":
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,
"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
CHAINED_JOBS_FILE = Path(__file__).parent / "chained-jobs.json"
def _has_chained(job_id):
try:
if CHAINED_JOBS_FILE.exists():
import json as _j
with open(CHAINED_JOBS_FILE) as f:
return job_id in _j.load(f)
except: pass
return False
def _mark_chained(job_id):
try:
import json as _j
c = []
if CHAINED_JOBS_FILE.exists():
with open(CHAINED_JOBS_FILE) as f: c = _j.load(f)
if job_id not in c:
c.append(job_id)
c = c[-1000:]
with open(CHAINED_JOBS_FILE, "w") as f: _j.dump(c, f)
except: pass
def trigger_chain_next(job_id, result_text, success=True):
"""If the completed job has on_success, on_failure, or chain_next, dispatch downstream."""
if _has_chained(job_id):
return
_mark_chained(job_id)
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
# 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
try:
with open(job_file, "r", encoding="utf-8") as f:
cfg = json.load(f)
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)
env = os.environ.copy()
env["CHAIN_PREV_JOB_ID"] = job_id
env["CHAIN_PREV_RESULT"] = result_text[:1000]
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)])
subprocess.Popen(
cmd,
env=env,
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
)
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")
except Exception:
pass
def clear_matching_followups(followups, agent, thread_id, mid, text, dry_run=False,
job_id=None, verb=None):
"""Resolve follow-up records if an assistant message is detected in the thread.
Matches on thread identity (thread_uuid or target='main') OR on job_id
(from a [RESULT <job_id>] reply). The job_id path works regardless of
thread_uuid or target, fixing ghost followups with null thread_uuid.
Also matches a main-chat reply when the sweeper recorded
final_nudge_target='main' (final nudge routed to main chat).
When verb is given ([ACK|CLAIM|RESULT|DECLINE|NO-ACTION <job_id>]),
matching is job_id-scoped for non-RESULT verbs (a verb marker names the
digest it answers, so it must not touch unrelated pending followups that
merely share the thread); RESULT keeps the historical thread-or-job
matching. Every verb match records outcome=<verb>. ACK/CLAIM set status
'acknowledged' (nudge-suppressed, NOT closed) instead of 'resolved', and
also match already-'acknowledged' records so an ACK -> RESULT lifecycle
closes correctly.
"""
if not followups:
return
modified = False
for f_id, f_rec in followups.items():
# Verb replies can follow an ACK (ACK -> RESULT lifecycle), so verbs
# also match 'acknowledged' records; plain replies only match pending.
if verb:
if f_rec.get("status") not in ("pending", "acknowledged"):
continue
elif 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
elif f_rec.get("final_nudge_target") == "main" and thread_id == "main":
# C3: the sweeper routed the final nudge to main chat, so a
# main-chat reply resolves even when the followup target is a
# sidechat.
match_thread = True
# Match by job_id (from [RESULT <job_id>] or [VERB <job_id>]) --
# works regardless of thread_uuid or target. This is an ADDITIONAL
# path, not a replacement.
match_job = False
if job_id and f_rec.get("job_id") and f_rec.get("job_id") == job_id:
match_job = True
# Non-RESULT verbs are job-scoped: they must not acknowledge/resolve
# unrelated pending followups that merely share the thread. RESULT
# keeps the historical thread-or-job matching.
if verb and verb != "RESULT" and not match_job:
continue
if match_thread or match_job:
if verb:
f_rec["outcome"] = verb
if verb in ("ACK", "CLAIM"):
# Acknowledged: sweeper nudges stop (status != pending), but
# the digest is NOT closed until a closing verb arrives.
f_rec["status"] = "acknowledged"
f_rec["acknowledged_at"] = utcnow()
f_rec["acknowledged_by_mid"] = mid
else:
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:
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
)
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
# 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:
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
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()
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()