1608 lines
65 KiB
Python
Executable File
1608 lines
65 KiB
Python
Executable File
def is_title_noise(title):
|
|
t = (title or "").lower().strip()
|
|
return bool(re.search(r"generate.*(chat|session)?.*title", t))
|
|
|
|
#!/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 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
|
|
|
|
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()
|
|
|
|
|
|
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",
|
|
}
|
|
|
|
|
|
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")
|
|
return op, args
|
|
|
|
|
|
def parse_tool_calls(text):
|
|
"""
|
|
Extract structured tool/exec calls from assistant messages.
|
|
Supports:
|
|
1. [TOOL <op> <json_args>] or [EXEC <op> <json_args>]
|
|
2. ```box / ```tool / ```exec JSON blocks
|
|
"""
|
|
calls = []
|
|
for m in re.finditer(r"\[(?:TOOL|EXEC)\s+([a-zA-Z0-9_.-]+)(?:\s+(.*?))?\]", text or ""):
|
|
op = m.group(1).strip()
|
|
raw_args = (m.group(2) or "").strip()
|
|
args = {}
|
|
if raw_args:
|
|
try:
|
|
args = json.loads(raw_args)
|
|
except Exception:
|
|
args = {"raw": raw_args}
|
|
calls.append((op, args))
|
|
|
|
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()
|
|
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)."
|
|
|
|
# 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)
|
|
|
|
|
|
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, or recent threads (< 3h).
|
|
"""
|
|
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)
|
|
|
|
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
|
|
|
|
if not (is_perm or is_pending or is_recent):
|
|
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
|
|
|
|
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":
|
|
markers = list(iter_result_markers(text))
|
|
verbs = list(iter_verb_markers(text))
|
|
|
|
# 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}"
|
|
tool_hint = (
|
|
f"[Runtime Context: {thread_url}]\n"
|
|
f"Tools available: [TOOL <op> <args>] or curl -sk -X POST https://exec.muse-dev.online/exec\n"
|
|
f" • [TOOL followup.create {{\"in_m\": 5, \"prompt\": \"<reminder>\"}}]\n"
|
|
f" • [TOOL swarm.spawn {{\"count\": 1, \"task\": \"<task>\"}}]\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:
|
|
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)
|
|
# 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:
|
|
archive_ephemeral_thread(agent, thread_id, job_id=job_id)
|
|
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)
|
|
|
|
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):
|
|
"""
|
|
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 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")
|
|
|
|
# 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.
|
|
SWARM_WORKER_POOL = ["dev", "def", "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 available auxiliary worker nodes (dev, def, muse).
|
|
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
|
|
prompt = (
|
|
f"[JOB {sid}/{slot_idx}] Task for swarm slot {slot_idx}:\n"
|
|
f"{task_text}\n\n"
|
|
f"Reply with [RESULT {sid}/{slot_idx}] OK <summary> or FAIL <reason>."
|
|
)
|
|
|
|
# Verified dispatch: confirm the send landed BEFORE
|
|
# marking the slot running. Fire-and-forget used to
|
|
# freeze slots from birth -- the slot read "running"
|
|
# while no worker was ever notified (2026-10-05: dev's
|
|
# wedged WireGuard data path silently killed every
|
|
# gateway dispatch).
|
|
dispatch_ok = False
|
|
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 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
|
|
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 slots or matching 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):
|
|
"""
|
|
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.
|
|
"""
|
|
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"}]'
|
|
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 executable tool calls now: [TOOL <op> <args>] 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.
|
|
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}
|
|
|
|
# 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()
|