Files
box/bin/response-harvester.py
T

2046 lines
85 KiB
Python
Executable File

#!/usr/bin/env python3
def is_title_noise(title):
t = (title or "").lower().strip()
return bool(re.search(r"generate.*(chat|session)?.*title", t))
"""
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 concurrent.futures
import hashlib
import json
import os
import re
import ssl
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")
SWARM_FILE = NETVM_ROOT / "swarms.json"
NUDGE_TRACKER_FILE = NETVM_ROOT / "conversation-nudge-tracker.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
try:
import muse_hybrid
HAS_MUSE_HYBRID = True
except ImportError:
HAS_MUSE_HYBRID = False
try:
import prompt_envelope
HAS_PROMPT_ENVELOPE = True
except ImportError:
HAS_PROMPT_ENVELOPE = False
VALID_AGENTS = ["muse", "pip", "646", "opm", "dev", "def"]
DEFAULT_PORTS = {"muse": 9410, "pip": 9420, "646": 9430, "opm": 9440, "def": 9450, "dev": 9455}
try:
import lookup_engine
HAS_LOOKUP_ENGINE = True
except ImportError:
HAS_LOOKUP_ENGINE = False
# 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.
_LOCAL_RESULT_RE = re.compile(r"\[RESULT\s+([A-Za-z0-9_/-]+)\]\s*(.*?)(?=\[RESULT\s|\Z)", re.S)
RESULT_RE = lookup_engine.get_result_regex() if HAS_LOOKUP_ENGINE else _LOCAL_RESULT_RE
def iter_result_markers(text):
"""Yield (job_id, result_text) for every [RESULT <job_id>] marker in text."""
current_re = lookup_engine.get_result_regex() if HAS_LOOKUP_ENGINE else RESULT_RE
for m in current_re.finditer(text or ""):
gd = m.groupdict()
if "summary" in gd:
# Engine shape: [RESULT <id>] [STATUS] <summary>. The
# status word is optional (None for bare markers); keep
# it when present so FAIL/ERROR still trips failure
# detection downstream.
status = (m.group("status") or "").strip()
summary = (m.group("summary") or "").strip()
result_text = f"{status} {summary}".strip() if status else summary
yield m.group("job_id").strip(), result_text
else:
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.
_LOCAL_VERB_RE = re.compile(r"\[(ACK|CLAIM|RESULT|DECLINE|NO-ACTION)\s+([A-Za-z0-9_/-]+)\]")
VERB_RE = lookup_engine.get_verb_regex() if HAS_LOOKUP_ENGINE else _LOCAL_VERB_RE
def iter_verb_markers(text):
"""Yield (verb, job_id) for every [VERB <job_id>] marker in text."""
current_re = lookup_engine.get_verb_regex() if HAS_LOOKUP_ENGINE else VERB_RE
for m in current_re.finditer(text or ""):
yield m.group(1), m.group(2).strip()
REFUSAL_PATTERNS = [
r"\bnot doing that\b",
r"\brefuse to\b",
r"\bdecline to\b",
r"\bdeclining\b",
r"\bwill not perform\b",
r"\bnot emitting the \[?result\]?\b",
r"\bnot sending result\b",
r"\bproduction change with recursion risk\b",
]
def detect_explicit_refusal(text):
"""Detect plain-language refusal/decline from an assistant in a task thread."""
t = (text or "").lower()
return any(re.search(pat, t) for pat in REFUSAL_PATTERNS)
# Muse-native tool names (docs.muse-dev.online/capabilities.html) -> box exec ops.
NATIVE_ALIASES = {
"cron.create": "followup.create",
"cron.runonce": "followup.create",
"cron.schedule": "followup.create",
"subagent.spawn": "swarm.spawn",
"subagents.spawn": "swarm.spawn",
"subagent.list": "swarm.list",
"subagent.status": "swarm.status",
"dm": "dm.send",
"message.send": "dm.send",
"box": "box.exec",
"box.run": "box.exec",
"tools": "tools.list",
"tools.list": "tools.list",
}
def normalize_native_call(op, args):
"""Bridge native Muse tool names/args to the box ops the exec daemon allows."""
op = NATIVE_ALIASES.get(op, op)
if not isinstance(args, dict):
return op, args
args = dict(args)
if op == "followup.create":
for k in ("kind", "mode", "schedule"): # native cron fields the box op does not take
args.pop(k, None)
if "in_m" not in args:
for k in ("minutes", "in_minutes", "delay_m", "runonce_m"):
if k in args:
args["in_m"] = args.pop(k)
break
if "prompt" not in args:
for k in ("body", "task", "message", "text"):
if k in args:
args["prompt"] = args.pop(k)
break
elif op == "swarm.spawn":
if "task" not in args:
for k in ("prompt", "body", "instructions"):
if k in args:
args["task"] = args.pop(k)
break
if "count" not in args and "n" in args:
args["count"] = args.pop("n")
elif op == "dm.send":
if "message" not in args:
for k in ("text", "body", "content", "msg"):
if k in args:
args["message"] = args.pop(k)
break
if "target" not in args:
for k in ("thread", "sidechat", "channel"):
if k in args:
args["target"] = args.pop(k)
break
elif op == "box.exec":
if "action" not in args:
for k in ("cmd", "verb", "command", "run"):
if k in args:
args["action"] = args.pop(k)
break
elif op.startswith("tmux."):
if "session" not in args:
for k in ("name", "target", "s"):
if k in args:
args["session"] = args.pop(k)
break
if op == "tmux.send" and "keys" not in args:
for k in ("command", "cmd", "input", "text"):
if k in args:
args["keys"] = args.pop(k)
break
return op, args
_TOOL_OPEN_RE = re.compile(r"\[(TOOL|EXEC)\s+([a-zA-Z0-9_.-]+)\s*")
_DM_OPEN_RE = re.compile(r"\[DM\s+")
def _extract_balanced_json(s, i):
"""Extract one JSON object starting at s[i] == '{' (brace-aware, string-aware).
Returns (obj, end_index) with end_index just past the closing brace,
or (None, i) when no balanced object is present. Unlike a first-']'
regex this tolerates ']' (and nested objects/arrays) inside args.
"""
if i >= len(s) or s[i] != "{":
return None, i
depth = 0
in_str = False
esc = False
for j in range(i, len(s)):
c = s[j]
if in_str:
if esc:
esc = False
elif c == "\\":
esc = True
elif c == '"':
in_str = False
elif c == '"':
in_str = True
elif c == "{":
depth += 1
elif c == "}":
depth -= 1
if depth == 0:
try:
return json.loads(s[i:j + 1]), j + 1
except Exception:
return None, i
return None, i
def _scan_bracket_calls(text):
"""Yield (op, args) for [TOOL op {...}] / [EXEC op {...}] / [DM {...}].
JSON args are extracted with balanced-brace scanning so ']' inside
strings, arrays, or nested objects no longer truncates the call.
Non-JSON tails keep the legacy first-']' behavior (raw passthrough).
"""
out = []
spans = []
for m in _TOOL_OPEN_RE.finditer(text or ""):
spans.append((m.start(), "tool", m.group(2).strip(), m.end()))
for m in _DM_OPEN_RE.finditer(text or ""):
spans.append((m.start(), "dm", "dm.send", m.end()))
spans.sort()
for _, kind, op, pos in spans:
if pos < len(text) and text[pos] == "{":
args, _ = _extract_balanced_json(text, pos)
if args is None:
continue
if not isinstance(args, dict):
args = {"raw": args}
elif kind == "dm":
continue # [DM ...] requires a JSON object; skip bare forms
else:
end = text.find("]", pos)
if end == -1:
continue
raw_args = text[pos:end].strip()
if not raw_args:
args = {}
else:
try:
args = json.loads(raw_args)
if not isinstance(args, dict):
args = {"raw": args}
except Exception:
args = {"raw": raw_args}
out.append((op, args))
return out
_PROOF_EVIDENCE_RE = re.compile(
r"sw-\d{8}-\d{6}-[0-9a-f]{4}"
r"|[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}"
r"|/(?:[\w.-]+/)+[\w.-]+"
r"|\b(?:swarm|timer|cron|job|thread|sidechat|slot)[-_ ]?(?:id|name|uuid)?\s*[:=]"
r"|\b\d+/\d+\s*(?:slots?|checks?|workers?)",
re.IGNORECASE)
_UUID_RE = re.compile(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}")
def result_has_evidence(result_text):
"""True when a RESULT verdict carries checkable artifacts (IDs, paths, counts)."""
return bool(_PROOF_EVIDENCE_RE.search(result_text or ""))
# Automated in-thread proof requests disabled per fleet governance decision (2026-10-09)
PROOF_REQUESTS_ENABLED = False
def maybe_request_proof(agent, thread_id, job_id, result_text, dry_run=False):
"""Ask for checkable evidence when a success RESULT has none.
Disabled by default per fleet decision 2026-10-09: automated in-thread proof
challenges trigger adversarial rejection loops and waste agent quota.
"""
if not PROOF_REQUESTS_ENABLED:
return False
"""Ask for checkable evidence when a success RESULT has none.
One-shot per (thread, job) via the nudge tracker. Returns True when a
proof followup was scheduled.
"""
if dry_run or not thread_id or not _UUID_RE.fullmatch(thread_id.lower()):
return False
if result_has_evidence(result_text):
return False
tracker = load_json_file(NUDGE_TRACKER_FILE)
rec = tracker.get(thread_id, {})
done = rec.get("proof_jobs", [])
if job_id in done:
return False
ok, res = execute_agent_tool(agent, "followup.create", {
"agent": agent,
"in_m": 30,
"thread": thread_id,
"prompt": (
f"[PROOF] Your [RESULT {job_id}] has no checkable evidence. "
f"Reply in this thread with the swarm/timer IDs, paths, or command output "
f"that prove the outcome — or say what is still missing."),
})
if ok:
rec["proof_jobs"] = (done + [job_id])[-50:]
tracker[thread_id] = rec
save_json_file(NUDGE_TRACKER_FILE, tracker)
append_jsonl(JOB_LOG, {
"ts": utcnow(),
"type": "proof_requested",
"job_id": job_id,
"agent": agent,
"thread_id": thread_id,
})
return True
sys.stderr.write(f"warning: proof followup failed for {job_id}: {res}\n")
return False
def parse_tool_calls(text):
"""
Extract structured tool/exec calls from assistant messages.
Supports:
1. [TOOL <op> <json_args>] or [EXEC <op> <json_args>] (JSON may nest)
2. [DM <json_args>] shorthand for dm.send
3. ```box / ```tool / ```exec JSON blocks
"""
calls = []
calls.extend(_scan_bracket_calls(text))
for m in re.finditer(r"```(?:box|tool|exec)\s*\n(.*?)```", text or "", re.DOTALL):
block = m.group(1).strip()
try:
d = json.loads(block)
if isinstance(d, dict) and "op" in d:
calls.append((d["op"], d.get("args", {})))
except Exception:
pass
# Extract curl command tool calls (e.g. curl ... -d '{"op": "...", "args": ...}')
for m in re.finditer(r"curl\s+.*?https?://(?:box\.muse-dev\.online|100\.123\.153\.75:8444|exec\.muse-dev\.online)/exec.*?-[dD]\s+(['\"])(.*?)\1", text or "", re.DOTALL):
raw_body = m.group(2)
try:
d = json.loads(raw_body)
if isinstance(d, dict) and "op" in d:
calls.append((d["op"], d.get("args", {})))
except Exception:
pass
# De-duplicate identical calls within one message (envelope repeats the spawn
# call at top and bottom; an agent echoing both must not run it twice).
seen, unique = set(), []
for op, args in [normalize_native_call(o, a) for o, a in calls]:
key = (op, json.dumps(args, sort_keys=True, default=str))
if key in seen:
continue
seen.add(key)
unique.append((op, args))
return unique
def execute_agent_tool(agent, op, args, timeout=30):
"""
Execute tool call via local exec-constrained HTTPS daemon over Tailscale.
Returns (ok, result_or_error_string).
"""
token = None
token_file = os.path.expanduser(f"~/.exec-tokens/{agent}")
if not os.path.exists(token_file):
token_file = os.path.expanduser("~/.exec-server-token")
if os.path.exists(token_file):
try:
with open(token_file) as f:
token = f.read().strip()
except Exception:
pass
url = "https://100.123.153.75:8444/exec"
payload = json.dumps({"op": op, "args": args}).encode("utf-8")
headers = {
"Content-Type": "application/json",
"User-Agent": f"Box-Harvester/{agent}",
}
if token:
headers["Authorization"] = f"Bearer {token}"
ctx = ssl.create_default_context()
ctx.check_hostname = False
ctx.verify_mode = ssl.CERT_NONE
req = urllib.request.Request(url, data=payload, headers=headers, method="POST")
try:
with urllib.request.urlopen(req, context=ctx, timeout=timeout) as resp:
data = resp.read().decode("utf-8")
res = json.loads(data)
if res.get("rc") == 0:
return True, res.get("stdout", "").strip() or "OK"
elif "stdout" in res:
return False, (res.get("stdout") or res.get("stderr") or "failed").strip()
elif "error" in res:
return False, res["error"]
return True, json.dumps(res)
except urllib.error.HTTPError as he:
try:
err_body = he.read().decode("utf-8")
return False, f"HTTP {he.code}: {err_body}"
except Exception:
return False, f"HTTP {he.code}"
except Exception as e:
return False, str(e)
def format_tool_result_for_chat(op, raw_output):
"""
Format tool output into a clean, human-readable chat message
instead of spewing a raw wall of unformatted JSON.
"""
if not raw_output:
return "OK"
try:
data = json.loads(raw_output)
except Exception:
s = str(raw_output).strip()
if op == "box.exec":
if len(s) > 900:
s = s[:900] + "\n…(truncated, refine the call for detail)"
return f"box result:\n```\n{s}\n```"
return s[:500] if len(s) > 500 else s
if op == "health.check" and isinstance(data, dict) and "fleet" in data:
nodes = data.get("fleet", [])
all_ok = all(n.get("cdp_ok") and n.get("proc_alive") for n in nodes)
up_count = sum(1 for n in nodes if n.get("cdp_ok") and n.get("proc_alive"))
lines = [f"Fleet Health: {'ALL GREEN' if all_ok else 'DEGRADED'} ({up_count}/{len(nodes)} nodes online)"]
for n in nodes:
st = "OK" if n.get("cdp_ok") and n.get("proc_alive") else "FAIL"
lat = n.get("latency_ms", 0)
lines.append(f" • {n.get('node')}: {st} ({lat}ms)")
return "\n".join(lines)
if op == "cron.runs" and isinstance(data, dict) and "jobs" in data:
jobs = data.get("jobs", [])
return f"{len(jobs)} scheduled jobs configured (e.g. {', '.join(j.get('name') for j in jobs[:6])})"
if op == "cron.status" and isinstance(data, dict):
return f"Cron '{data.get('name')}': {'ACTIVE' if data.get('active') else 'INACTIVE'} (last result: {data.get('last_result', 'unknown')})"
if op == "vars.list" and isinstance(data, dict) and "variables" in data:
vars_dict = data.get("variables", {})
sample = ", ".join(f"{k}={v}" for k, v in list(vars_dict.items())[:5])
return f"{len(vars_dict)} variables: {sample}..."
if op == "vars.get" and isinstance(data, dict):
return f"{data.get('name')} = {data.get('value')}"
if op == "files.read" and isinstance(data, dict):
if not data.get("ok"):
return f"Failed to read file: {data.get('error')}"
path = data.get("path")
lines = data.get("lines")
trunc = " (first " + str(data.get("displayed_lines")) + " lines)" if data.get("truncated") else ""
content = data.get("content", "").strip()
return f"File `{path}` ({lines} lines total){trunc}:\n```\n{content}\n```"
if op == "files.write" and isinstance(data, dict):
if not data.get("ok"):
return f"Failed to write file: {data.get('error')}"
return f"File `{data.get('path')}` written successfully ({data.get('bytes_written')} bytes, {data.get('lines')} lines)."
if op == "web.fetch" and isinstance(data, dict):
if not data.get("ok"):
return f"Failed to fetch {data.get('url')}: {data.get('error')}"
status = data.get("status")
url = data.get("url")
trunc = " (truncated to 4KB)" if data.get("truncated") else ""
return f"Web Fetch `{url}` (HTTP {status}){trunc}:\n```\n{data.get('body', '').strip()[:800]}\n```"
if op == "service.status" and isinstance(data, dict):
if not data.get("ok"):
return f"Service check failed: {data.get('error')}"
st = "ACTIVE" if data.get("active") else "INACTIVE"
return f"Service `{data.get('service')}` is {st}.\nStatus: {data.get('status_line')}"
if op in ("followup.create", "followup.schedule") and isinstance(data, dict):
if not data.get("ok"):
return f"Follow-up scheduling failed: {data.get('error')}"
sec = data.get("in_seconds", 60)
ag = data.get("agent", "")
return f"Follow-up timer scheduled for {ag} in {sec}s ({round(sec/60, 1)}m)."
if op == "swarm.spawn" and isinstance(data, dict):
if not data.get("ok"):
return f"Swarm spawn failed: {data.get('error')}"
return f"Swarm `{data.get('swarm_id')}` spawned with {data.get('count')} worker slots (status: {data.get('status')})."
if op == "swarm.status" and isinstance(data, dict):
if not data.get("ok"):
return f"Swarm status check failed: {data.get('error')}"
sw = data.get("swarm", {})
return f"Swarm `{sw.get('swarm_id')}`: {sw.get('status')} ({sw.get('done', 0)}/{sw.get('count', 0)} slots completed)."
if op == "tools.list" and isinstance(data, dict):
ops = data.get("ops", [])
if not ops:
return "No tools registered."
ro = [o["op"] for o in ops if not o.get("side_effecting")]
se = [o["op"] for o in ops if o.get("side_effecting")]
lines = [f"{len(ops)} tools available via [TOOL <op> <args>]."]
lines.append("read-only: " + (", ".join(ro) if ro else "none"))
if se:
lines.append("side-effecting: " + ", ".join(se))
return "\n".join(lines)
if op == "box.exec" and isinstance(data, dict):
if data.get("ok") is False:
return f"box call failed: {data.get('error', 'unknown error')}"
out = json.dumps(data)
if len(out) > 900:
out = out[:900] + "\n…(truncated, refine the call for detail)"
return f"box result:\n```\n{out}\n```"
if op == "flow.start" and isinstance(data, dict):
if not data.get("ok"):
return f"Flow start failed: {data.get('error')}"
return f"Flow `{data.get('flow_id')}` started in pane `{data.get('session')}` (status: {data.get('status')})."
if op == "flow.read" and isinstance(data, dict):
if not data.get("ok"):
return f"Flow read failed: {data.get('error')}"
st = data.get("status", "unknown")
ec = data.get("exit_code")
ec_str = f" (exit_code: {ec})" if ec is not None else ""
pm = data.get("prompt_match")
prompt_str = f"\nPrompt waiting: {pm.get('text', pm)}" if pm else ""
delta = data.get("delta", "").strip()
trunc = f" (last {data.get('lines_read')} lines)" if data.get("truncated") else ""
body = f"\n```\n{delta}\n```" if delta else " (no new output)"
return f"Flow `{data.get('flow_id')}` [{st}]{ec_str}{prompt_str}{trunc}:{body}"
if op == "flow.send" and isinstance(data, dict):
if not data.get("ok"):
return f"Flow send failed: {data.get('error')}"
kind = "command" if data.get("is_command") else "keys"
return f"Flow `{data.get('flow_id')}` sent {kind}: `{data.get('sent')}` (status: {data.get('status')})."
if op == "flow.list" and isinstance(data, dict):
flows = data.get("flows", [])
if not flows:
return "No active flows."
lines = [f"{len(flows)} flows:"]
for f in flows[:8]:
lines.append(f" • {f.get('flow_id')} [{f.get('status')}]: {f.get('session')} (cmd: {str(f.get('command', 'bash'))[:30]})")
return "\n".join(lines)
if op == "flow.stop" and isinstance(data, dict):
return f"Flow `{data.get('flow_id')}` stopped."
# General fallback: compact JSON capped to 400 chars
s = json.dumps(data)
return s[:400] + "..." if len(s) > 400 else s
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)
RECENCY_WINDOW_SEC = 10800
_JOB_ID_RE = re.compile(r"^(.+)-(\d{8})-(\d{6})-([0-9a-f]{8})$")
def dispatched_families_since(job_log_path, window_sec=RECENCY_WINDOW_SEC,
now=None):
"""Job families dispatched inside the window.
Scans job-log.jsonl for job_sent/job_dispatched events newer than
``window_sec`` and returns their family names (the job id minus the
trailing -YYYYMMDD-HHMMSS-<hash> run suffix). Missing, unreadable,
or malformed input yields an empty set, never an exception.
"""
now = now or datetime.now(timezone.utc)
cutoff = now.timestamp() - window_sec
fams = set()
try:
handle = open(job_log_path, "r", encoding="utf-8")
except OSError:
return fams
with handle:
for line in handle:
line = line.strip()
if not line:
continue
try:
event = json.loads(line)
except Exception:
continue
if event.get("type") not in ("job_sent", "job_dispatched"):
continue
try:
ts = datetime.fromisoformat(
str(event.get("ts")).replace("Z", "+00:00")).timestamp()
except Exception:
continue
if ts < cutoff:
continue
match = _JOB_ID_RE.match(str(event.get("job_id") or ""))
if match:
fams.add(match.group(1))
return fams
def get_monitored_threads(target_agent=None):
"""
Build dict of threads to monitor per agent:
{ agent: [ {"id": "<uuid>", "name": "<alias>"} ] }
Filters to permanent channels, threads with pending followups,
recently created threads (< 3h), or threads whose job family was
dispatched recently (< 3h) so old persistent sidechats that still
receive prompts stay monitored.
"""
agents = [target_agent] if target_agent else VALID_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)
recently_dispatched = dispatched_families_since(JOB_LOG, now=now)
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
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 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 already archived, exclude from active monitoring unless pending followup
if isinstance(val, dict) and val.get("archived") and not is_pending:
continue
is_dispatched = key in recently_dispatched
if not (is_perm or is_pending or is_recent or is_dispatched):
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
# Also monitor active swarm slot subagents from swarms.json
try:
if SWARM_FILE.exists():
s_data = json.loads(SWARM_FILE.read_text(encoding="utf-8"))
for sid, s_info in s_data.items():
if sid.startswith("_") or not isinstance(s_info, dict):
continue
if s_info.get("status") not in ("running", "pending"):
continue
for slot in s_info.get("slots", []):
if slot.get("status") == "running":
sub_id = slot.get("subagent_session_id") or slot.get("thread_uuid") or slot.get("sidechat_id")
ag = slot.get("agent_id")
if sub_id and ag and ag in threads_by_agent:
existing = [t["id"] for t in threads_by_agent[ag]]
if sub_id not in existing:
threads_by_agent[ag].append({"id": sub_id, "name": f"swarm-{sid[:8]}-s{slot.get('slot')}"})
except Exception as se:
sys.stderr.write(f"warning: failed to add swarm threads to monitor list: {se}\n")
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 process_messages(raw_messages, agent, thread_id, thread_name, last_wm, followups, dry_run=False, feed=None):
"""
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)
"""
if not raw_messages:
return [], last_wm, 0
new_messages = []
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"]
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:
# Watermark not found in loaded window
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": thread_id,
"thread_name": thread_name,
"msg_id": mid,
"author": author,
"text": text,
"source_ts": msg_ts,
}
if feed:
record["feed"] = feed
if not dry_run:
append_jsonl(CHAT_HISTORY_LOG, record)
if author == "assistant":
try:
markers = list(iter_result_markers(text))
verbs = list(iter_verb_markers(text))
except Exception as e:
# One poison message must not wedge the batch: without
# this, the same crash repeats every cycle, the
# watermark never advances past it, and the thread's
# followups nag to escalation despite answered work.
sys.stderr.write(
"warning: marker extraction failed, treating as "
f"plain reply: {e}\n")
markers, verbs = [], []
# Synthesize [RESULT <job-id>] DECLINE if assistant explicitly refuses the task in plain text
if not markers and not verbs and detect_explicit_refusal(text):
matching_jid = None
for f_id, f_rec in followups.items():
if f_rec.get("status") in ("pending", "acknowledged") and (f_rec.get("thread_uuid") == thread_id or f_rec.get("recipient") == agent):
matching_jid = f_rec.get("job_id") or f_id
break
if not matching_jid and JOB_SIDECHATS_FILE.exists():
try:
sc_state = json.loads(JOB_SIDECHATS_FILE.read_text(encoding="utf-8"))
for k, v in sc_state.items():
if isinstance(v, dict) and v.get("thread_uuid") == thread_id:
matching_jid = k
break
except Exception:
pass
if matching_jid:
refusal_snippet = f"DECLINE: {text.strip()[:180]}"
print(f"[{agent}] Synthesizing [RESULT {matching_jid}] DECLINE from explicit assistant refusal")
markers.append((matching_jid, refusal_snippet))
# Check for structured tool / command calls from the agent
tool_calls = parse_tool_calls(text)
if tool_calls and not dry_run:
for op, args in tool_calls:
if isinstance(args, dict) and "agent" not in args:
args["agent"] = agent
print(f"[{agent}] Executing tool '{op}' from message {mid[:8]} in thread {thread_name or thread_id[:8]}")
t_ok, t_res = execute_agent_tool(agent, op, args)
tool_record = {
"ts": utcnow(),
"type": "tool_exec",
"agent": agent,
"thread_id": thread_id,
"op": op,
"args": args,
"success": t_ok,
"result_snippet": str(t_res)[:300],
"msg_id": mid,
}
append_jsonl(JOB_LOG, tool_record)
# Post clean tool output back to the originating thread via fast headless gateway
if thread_id and re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", thread_id.lower()):
try:
import muse_hybrid
thread_url = f"https://box.muse-dev.online/thread/{thread_id}"
if op.startswith("flow."):
flow_id = t_args.get("flow_id", "<flow_id>") if isinstance(t_args, dict) else "<flow_id>"
tool_hint = (
f"[Flow Directive: advance with [TOOL flow.send {{\"flow_id\": \"{flow_id}\", \"command\": \"...\"}}]"
f" | read with [TOOL flow.read {{\"flow_id\": \"{flow_id}\"}}]"
f" | close with [RESULT <job_id>] OK]"
)
else:
tool_hint = (
f"[Runtime Context: {thread_url}]\n"
f"Tools: EMIT one [TOOL <op> <args>] line per action (you do not run it;"
f" the runtime executes it and replies here). curl -sk -X POST"
f" https://exec.muse-dev.online/exec works too.\n"
f" • [TOOL tools.list {{}}] — discover every op dynamically\n"
f" • [TOOL swarm.spawn {{\"count\": 1, \"task\": \"<task>\"}}] — spawn subagents\n"
f" • [DM {{\"to\": \"<agent>\", \"target\": \"<sidechat>\", \"message\": \"<text>\"}}] — send a DM\n"
f" • [TOOL box.exec {{\"action\": \"fleet-status\"}}] — call box (read-only actions)\n"
f" • [TOOL followup.create {{\"in_m\": 5, \"prompt\": \"<reminder>\"}}]\n"
f" • [TOOL health.check {{}}]\n\n"
f"[Directive: Take next action or close with [RESULT <job_id>] <summary>]"
)
if t_ok:
clean_msg = format_tool_result_for_chat(op, t_res)
resp_text = f"Tool result (`{op}`):\n{clean_msg}\n\n{tool_hint}"
else:
resp_text = f"Tool error (`{op}`): {t_res}\n\n{tool_hint}"
muse_hybrid.send_message(agent, resp_text, thread_id=thread_id, wait=0)
except Exception as te:
sys.stderr.write(f"warning: failed to post tool response back to thread: {te}\n")
if markers or verbs:
seen_jobs = set()
for job_id, result_text in markers:
if job_id in seen_jobs:
# Same verdict restated in one message: log once.
continue
seen_jobs.add(job_id)
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 result_text.startswith("DECLINE:"):
# Synthesized (or explicit) decline: still a
# non-success (no chaining), but the auditor
# buckets it as declined, not a failure.
job_record["outcome"] = "declined"
if not dry_run:
append_jsonl(JOB_LOG, job_record)
# Check if this is a swarm slot result: sw-YYYYMMDD-HHMMSS-xxxx/<slot>
if "/" in job_id and job_id.startswith("sw-"):
try:
s_sid, s_slot = job_id.split("/", 1)
s_proc = subprocess.Popen(
[sys.executable, str(BIN_DIR / "box-ctl.py"), "swarm-report", s_sid, s_slot],
stdin=subprocess.PIPE, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL
)
s_proc.communicate(input=json.dumps({"ok": not is_fail, "result": result_text}).encode("utf-8"))
except Exception as se:
sys.stderr.write(f"warning: failed to record swarm report: {se}\n")
else:
trigger_chain_next(job_id, result_text, success=not is_fail)
if not is_fail:
try:
maybe_request_proof(agent, thread_id, job_id, result_text,
dry_run=dry_run)
except Exception as pe:
sys.stderr.write(f"warning: proof check failed: {pe}\n")
archive_ephemeral_thread(agent, thread_id, job_id=job_id,
dry_run=dry_run)
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
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)
maybe_nudge_untagged_sidechat(agent, thread_id, thread_name, mid, text,
dry_run=dry_run, acted=bool(tool_calls))
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
)
def archive_ephemeral_thread(agent, thread_id, job_id=None, dry_run=False):
"""
If thread_id belongs to an ephemeral job or one-off check,
archive it via hybrid gateway and tag it as archived in job-sidechats.json.
"""
if dry_run:
return
if not thread_id or not re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", thread_id.lower()):
return
# Check if thread is marked as persistent in job-sidechats.json
try:
if JOB_SIDECHATS_FILE.exists():
sc_state = json.loads(JOB_SIDECHATS_FILE.read_text(encoding="utf-8"))
for key, val in sc_state.items():
if isinstance(val, dict) and val.get("thread_uuid") == thread_id:
if val.get("type") == "persistent":
return # Do not archive persistent coordinator channels
val["archived"] = True
val["archived_at"] = utcnow()
val["archived_by_job"] = job_id
JOB_SIDECHATS_FILE.write_text(json.dumps(sc_state, indent=2), encoding="utf-8")
except Exception as e:
sys.stderr.write(f"warning: failed to update job-sidechats state: {e}\n")
# Complete subagent tracker record if this thread corresponds to an ephemeral subagent session
try:
import subagent_tracker
subagent_tracker.complete_session(thread_id, note=f"Harvested verdict for job {job_id}")
except Exception:
pass
# Call muse-threads.py archive via subprocess (runs in node's netns)
try:
helper = BIN_DIR / "muse-threads.py"
subprocess.run(
[sys.executable, str(helper), "archive", "--agent", agent, "--thread", thread_id],
capture_output=True, text=True, timeout=15
)
print(f"[{agent}] Archived ephemeral thread {thread_id[:8]} upon completion of job {job_id or 'unknown'}")
except Exception as ae:
sys.stderr.write(f"warning: archive_ephemeral_thread failed: {ae}\n")
# NOTE 2026-10-05 (Fix Agent 2/5): dev/def removed from the pool. No worker
# agents exist on dev/def (no Meta sessions provisioned), so slots dispatched
# to them froze with null results. Re-add only after real dev/def workers exist.
# NOTE 2026-10-06 (split-brain dispatch fix): the removal above was comment-only;
# the code still listed dev/def, and the harvester's every-minute dispatch raced
# the pool daemon claiming slots as dev/def, which refuse on attribution grounds
# (dev FAILs fast, def freezes until the 60-min reaper). Pool now matches
# swarm_worker/daemon.py WORKER_POOL exactly: only muse, the single fully
# authenticated auxiliary worker.
SWARM_WORKER_POOL = ["muse"]
# Stuck-slot reaper: a slot that stays "running" with no result longer than
# this is treated as wedged (worker died / dispatch lost). Healthy slots
# complete in <5 min (observed p90 3.6 min over 26 done slots, 2026-10-05),
# so 60 min is conservative.
STUCK_SLOT_MINUTES = 60
# Verified dispatch: dm.py send runs synchronously and the slot is only
# marked running when the send is confirmed (rc 0 + "SENT" in stdout).
# Unverified sends stay pending for retry; after N attempts the slot fails
# loudly instead of freezing from birth.
DISPATCH_VERIFY_TIMEOUT = 120
DISPATCH_MAX_ATTEMPTS = 5
def _parse_ts(ts):
try:
return datetime.fromisoformat(str(ts).replace("Z", "+00:00"))
except Exception:
return None
def reconcile_and_dispatch_swarms(dry_run=False):
"""
Autonomous Swarm Orchestrator:
1. Scans swarms.json for pending slots.
2. Dynamically allocates the auxiliary worker pool (muse only; dev/def have
no provisioned worker agents and refuse on attribution grounds).
3. Provisions ephemeral sidechat per slot and dispatches the task with [RESULT <swarm_id>/<slot>].
4. Upon completion of all slots, sends completion summary DM to originating coordinator.
"""
if not SWARM_FILE.exists() or dry_run:
return
try:
swarms = json.loads(SWARM_FILE.read_text(encoding="utf-8"))
except Exception:
return
modified = False
now = utcnow()
# 0. Reap stuck slots BEFORE computing busy workers. A wedged worker
# would otherwise pin itself "busy" forever and the pool stalls: with
# all workers busy the fallback round-robin keeps feeding new slots to
# the same dead workers.
now_dt = _parse_ts(now)
if now_dt is not None:
for _sid, _swarm in swarms.items():
if _swarm.get("status") not in ("pending", "running"):
continue
for _s in _swarm.get("slots", []):
if _s.get("status") != "running" or _s.get("result") is not None:
continue
_upd = _parse_ts(_s.get("updated_ts", ""))
if _upd is None:
continue
_age_min = (now_dt - _upd).total_seconds() / 60
if _age_min > STUCK_SLOT_MINUTES:
_s["status"] = "failed"
_s["result"] = {
"ok": False,
"reaped": True,
"reason": "stuck: running with no result for %.0f min (limit %d)"
% (_age_min, STUCK_SLOT_MINUTES),
"worker": _s.get("agent_id"),
}
_s["updated_ts"] = now
modified = True
print("[swarm] Reaped stuck slot %d of %s (worker %s, silent %.0f min)"
% (_s["slot"], _sid, _s.get("agent_id"), _age_min))
# Determine busy workers from running slots
busy_workers = set()
for sid, swarm in swarms.items():
if swarm.get("status") in ("pending", "running"):
for slot in swarm.get("slots", []):
if slot.get("status") == "running" and slot.get("agent_id"):
busy_workers.add(slot["agent_id"])
# Process each swarm
for sid, swarm in swarms.items():
st = swarm.get("status")
creator = swarm.get("created_by") or "646"
if creator not in ("646", "pip", "opm", "muse", "dev", "def"):
creator = "646"
# 1. Allocate & dispatch pending slots
if st in ("pending", "running"):
slots = swarm.get("slots", [])
for s in slots:
if s.get("status") == "pending":
# Find first available worker
worker = None
for w in SWARM_WORKER_POOL:
if w not in busy_workers:
worker = w
break
if not worker:
# Fallback round-robin across worker pool if all are busy
worker = SWARM_WORKER_POOL[s["slot"] % len(SWARM_WORKER_POOL)]
slot_idx = s["slot"]
task_text = swarm.get("task", "execute subagent task")
slot_target = f"{sid}-s{slot_idx}"
# Prepare task directive prompt
slot_job_id = f"{sid}/{slot_idx}"
if HAS_PROMPT_ENVELOPE and hasattr(prompt_envelope, "wrap_subagent_task"):
prompt = prompt_envelope.wrap_subagent_task(slot_job_id, task_text)
else:
prompt = (
f"Operator assignment for swarm slot {slot_job_id}:\n\n"
f"Task:\n{task_text.strip()}\n\n"
f"Instructions:\n"
f"1. Carry out this task directly using your available tools.\n"
f"2. When finished, conclude your final response with your verdict line:\n"
f"[RESULT {slot_job_id}] OK: <one-line summary of actions and outcome>\n"
f"(or [RESULT {slot_job_id}] FAIL: <reason> if the task could not be completed)\n"
)
# Native Subagent Execution Bridge:
# Spawns an interactive child session directly on the target worker
# via muse_hybrid (Option 3), ensuring an active consumer actually executes
# the prompt and returns results.
dispatch_ok = False
subagent_sid = None
try:
import muse_hybrid
sub_res, sub_err = muse_hybrid.start_session(worker, title=f"swarm-{sid[:8]}-s{slot_idx}")
if not sub_err and sub_res and sub_res.get("session_id"):
subagent_sid = sub_res["session_id"]
send_res, send_err = muse_hybrid.send_message(worker, prompt, thread_id=subagent_sid, wait=0)
if not send_err:
dispatch_ok = True
try:
import subagent_tracker
subagent_tracker.register_session(creator, subagent_sid, title=f"{sid}-s{slot_idx}", prompt=prompt)
except Exception:
pass
print(f"[swarm] Dispatched slot {slot_idx} of {sid} to native subagent session {subagent_sid} on {worker}")
if not dispatch_ok:
sys.stderr.write(f"warning: native subagent dispatch failed for {sid} slot {slot_idx} -> {worker} (err: {sub_err or send_err})\n")
except Exception as se:
sys.stderr.write(f"warning: muse_hybrid subagent dispatch exception: {se}\n")
# Fallback to dm.py send if native session creation fails
if not dispatch_ok:
try:
dm_cmd = [
sys.executable, str(BIN_DIR / "dm.py"), "send",
"--agent", creator,
"--to", worker,
"--target", slot_target,
prompt,
]
proc = subprocess.run(
dm_cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE,
text=True, timeout=DISPATCH_VERIFY_TIMEOUT)
_dout = proc.stdout or ""
dispatch_ok = (proc.returncode == 0 and "SENT" in _dout)
if not dispatch_ok:
sys.stderr.write(
"warning: swarm fallback dispatch unverified %s slot %d -> %s "
"(rc=%d out=%.100s err=%.100s)\n"
% (sid, slot_idx, worker, proc.returncode,
_dout, proc.stderr or ""))
except Exception as de:
sys.stderr.write(
"warning: swarm dispatch error %s slot %d -> %s: %s\n"
% (sid, slot_idx, worker, de))
if not dispatch_ok:
# Leave pending so the next harvester pass retries
# (possibly with a different worker). Fail loudly
# after N attempts instead of spinning forever.
attempts = s.get("dispatch_attempts", 0) + 1
s["dispatch_attempts"] = attempts
modified = True
if attempts >= DISPATCH_MAX_ATTEMPTS:
s["status"] = "failed"
s["result"] = {
"ok": False,
"reaped": True,
"reason": "dispatch failed %d times (send never verified)" % attempts,
"worker": worker,
}
s["updated_ts"] = now
print("[swarm] Dispatch failed %dx for slot %d of %s; marked failed"
% (attempts, slot_idx, sid))
else:
print("[swarm] Dispatch unverified for slot %d of %s "
"(attempt %d/%d); leaving pending for retry"
% (slot_idx, sid, attempts, DISPATCH_MAX_ATTEMPTS))
continue
# Attach worker to slot (only after verified dispatch)
s["agent_id"] = worker
s["status"] = "running"
s["updated_ts"] = now
if subagent_sid:
s["subagent_session_id"] = subagent_sid
busy_workers.add(worker)
modified = True
print(f"[swarm] Dispatched slot {slot_idx} of {sid} to {worker} in sidechat {slot_target}")
# Update swarm rollup status
counts = {"pending": 0, "running": 0, "done": 0, "failed": 0, "killed": 0}
for slot in slots:
counts[slot.get("status", "pending")] = counts.get(slot.get("status", "pending"), 0) + 1
if counts["pending"] + counts["running"] == 0:
swarm["status"] = "completed" if counts["failed"] == 0 else "partial"
swarm["updated_ts"] = now
modified = True
elif counts["running"] or counts["done"]:
swarm["status"] = "running"
swarm["updated_ts"] = now
modified = True
# 2. Check if newly completed/partial and notify coordinator
if swarm.get("status") in ("completed", "partial") and not swarm.get("notified_coordinator"):
swarm["notified_coordinator"] = True
modified = True
done_cnt = sum(1 for sl in swarm.get("slots", []) if sl.get("status") == "done")
total_cnt = len(swarm.get("slots", []))
summary_msg = (
f"[Swarm Report] Swarm {sid} ({swarm.get('label') or 'task'}) finished: "
f"{done_cnt}/{total_cnt} slots successful. "
f"Console: https://box.muse-dev.online/#dashboard"
)
try:
coord_dm = [
sys.executable, str(BIN_DIR / "dm.py"), "send",
"--agent", "box",
"--to", creator,
"--target", "main",
"--allow-main-chat",
summary_msg,
]
subprocess.Popen(coord_dm, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL)
print(f"[swarm] Delivered completion report for {sid} to coordinator {creator}")
except Exception as ce:
sys.stderr.write(f"warning: failed to notify swarm coordinator: {ce}\n")
if modified:
tmp_swarms = str(SWARM_FILE) + f".tmp.{os.getpid()}"
with open(tmp_swarms, "w", encoding="utf-8") as f:
json.dump(swarms, f, indent=2)
os.replace(tmp_swarms, SWARM_FILE)
def check_and_archive_terminal_swarms():
"""Check swarms.json and auto-archive ephemeral sidechats for completed or terminal swarms."""
if not SWARM_FILE.exists() or not JOB_SIDECHATS_FILE.exists():
return
try:
swarms_data = json.loads(SWARM_FILE.read_text(encoding="utf-8"))
sc_data = json.loads(JOB_SIDECHATS_FILE.read_text(encoding="utf-8"))
except Exception:
return
for sid, swarm in swarms_data.items():
st = swarm.get("status")
if st in ("completed", "partial", "killed"):
# Check slot subagent sessions
for slot in swarm.get("slots", []):
sub_id = slot.get("subagent_session_id")
ag = slot.get("agent_id")
if sub_id and ag:
archive_ephemeral_thread(ag, sub_id, job_id=sid)
# Check matching registered sidechats
for key, val in sc_data.items():
if isinstance(val, dict) and not val.get("archived") and val.get("type") != "persistent":
if sid in key or (swarm.get("label") and swarm.get("label") in key):
tu = val.get("thread_uuid")
ag = val.get("agent", "opm")
if tu:
archive_ephemeral_thread(ag, tu, job_id=sid)
def maybe_nudge_untagged_sidechat(agent, thread_id, thread_name, mid, text, dry_run=False,
acted=False):
"""
If an agent replies conversationally in a sidechat backed by a job or follow-up
without providing [RESULT <id>] or tool directives, deliver a terse 1-turn nudge footer.
When acted=True the agent DID emit directives but never closed: remind to close
with [RESULT] instead of rejecting the (good) action.
"""
if dry_run or not thread_id or thread_id == "main":
return
if not re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", thread_id.lower()):
return
# Look up active follow-ups or job association for this thread
followups = load_json_file(FOLLOWUPS_FILE)
matching_job_id = None
for f_id, f_rec in followups.items():
if f_rec.get("status") in ("pending", "acknowledged") and f_rec.get("thread_uuid") == thread_id:
matching_job_id = f_rec.get("job_id") or f_id
break
if not matching_job_id and JOB_SIDECHATS_FILE.exists():
try:
sc_state = json.loads(JOB_SIDECHATS_FILE.read_text(encoding="utf-8"))
for k, v in sc_state.items():
if isinstance(v, dict) and v.get("thread_uuid") == thread_id:
if k.startswith(("box-", "job-", "pipe-", "autonomy-pulse-")):
matching_job_id = k
break
except Exception:
pass
if not matching_job_id:
return
tracker = load_json_file(NUDGE_TRACKER_FILE)
rec = tracker.get(thread_id, {})
if rec.get("nudged_for_mid") == mid:
return
n_count = rec.get("nudge_count", 0)
# 1st turn: Strict rejection and tool invocation prompt
if n_count == 0:
thread_url = f"https://box.muse-dev.online/thread/{thread_id}"
try:
import prompt_envelope
_spawn = prompt_envelope.spawn_call(
matching_job_id, matching_job_id, prompt_envelope.pick_profile(matching_job_id))
except Exception:
_spawn = '[TOOL swarm.spawn {"count": 2, "task": "continue the job work"}]'
if acted:
nudge_msg = (
f"Action received — now close the loop: reply with [RESULT {matching_job_id}] <outcome>.\n"
f"Outcome needs checkable evidence (swarm/timer IDs, paths, or command output), not prose alone.\n"
f"Thread Console: {thread_url}"
)
else:
nudge_msg = (
f"{_spawn}\n"
f"[STRICT ENFORCEMENT: Conversational commentary is rejected. Work requires active execution.]\n"
f"Thread Console: {thread_url}\n"
f"EMIT tool calls verbatim in your reply — you do not run them yourself;"
f" the Box runtime on bl executes each directive and posts the result back here"
f" (works from containers with no box CLI). Or curl against https://exec.muse-dev.online/exec\n"
f"When all operations are finished, close strictly with [RESULT {matching_job_id}] <outcome>.\n"
f"{_spawn}"
)
try:
import muse_hybrid
print(f"[{agent}] Injecting 1-turn strict nudge into {thread_name or thread_id[:8]} for job {matching_job_id}")
muse_hybrid.send_message(agent, nudge_msg, thread_id=thread_id, wait=0)
tracker[thread_id] = {
"ts": utcnow(),
"job_id": matching_job_id,
"nudged_for_mid": mid,
"nudge_count": 1,
}
save_json_file(NUDGE_TRACKER_FILE, tracker)
except Exception as ne:
sys.stderr.write(f"warning: failed to deliver conversational nudge: {ne}\n")
# 2nd turn: Trigger emergency 5m escalation follow-up to opm/coordinator
elif n_count == 1 and not rec.get("escalated"):
print(f"[{agent}] Conversational stall persisted in {thread_name or thread_id[:8]} - scheduling emergency escalation to opm")
esc_prompt = f"[ESCALATION] Agent {agent} stalled in thread {thread_name or thread_id[:8]} on job {matching_job_id}. Review thread or reassign."
try:
cmd = [
"systemd-run", "--user", "--on-active=300s",
sys.executable, str(BIN_DIR / "box-ctl.py"), "notify", "opm", esc_prompt, "--sender", agent
]
subprocess.run(cmd, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL)
tracker[thread_id]["escalated"] = True
tracker[thread_id]["escalated_ts"] = utcnow()
save_json_file(NUDGE_TRACKER_FILE, tracker)
except Exception as ee:
sys.stderr.write(f"warning: failed to schedule emergency escalation: {ee}\n")
# 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. Also matches the followup's own key
# (dm_id): agents quote the DM id from nudge text ([RESULT
# <dm_id>]), which differs from job_id on DM-ordered followups
# (observed live: [RESULT f4293153] vs job ml-muse-*).
match_job = False
if job_id and f_rec.get("job_id") and f_rec.get("job_id") == job_id:
match_job = True
elif job_id and job_id == f_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}
# 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:
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)
reconcile_and_dispatch_swarms()
check_and_archive_terminal_swarms()
# 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()