feat(work): add auto-heal engine, chat filtering, and CDP revival to box work
This commit is contained in:
+112
-13
@@ -149,7 +149,7 @@ def get_recent_done_tasks(limit=5):
|
||||
done.append({"name": f.name, "mtime": mtime.strftime("%H:%M:%SZ")})
|
||||
return done
|
||||
|
||||
def get_recent_chat_events(limit=5):
|
||||
def get_recent_chat_events(limit=5, agent=None):
|
||||
events = []
|
||||
if CHAT_LOG.exists():
|
||||
try:
|
||||
@@ -160,6 +160,8 @@ def get_recent_chat_events(limit=5):
|
||||
continue
|
||||
try:
|
||||
ev = json.loads(line)
|
||||
if agent and ev.get("agent") != agent:
|
||||
continue
|
||||
events.append(ev)
|
||||
if len(events) >= limit:
|
||||
break
|
||||
@@ -377,9 +379,100 @@ def cmd_check(args):
|
||||
print(f" • Hatch: {res['hatch']['details']}")
|
||||
print(f" • Restore: {res['restore']['details']}")
|
||||
print(f" • Git: {res['git']['details']}")
|
||||
for r in res["reasons"]:
|
||||
print(f" {c_yellow('!')} {r}")
|
||||
print()
|
||||
|
||||
def heal_agent(agent_name: str) -> dict:
|
||||
"""Automated remediation for an agent failing pre-flight health checks."""
|
||||
worker = next((w for w in WORKERS if w["name"] == agent_name), None)
|
||||
actions = []
|
||||
unresolved = []
|
||||
|
||||
if not worker and agent_name != "super":
|
||||
return {
|
||||
"agent": agent_name,
|
||||
"healed": False,
|
||||
"actions": [],
|
||||
"unresolved": [f"Unknown worker '{agent_name}'"]
|
||||
}
|
||||
|
||||
port = worker["port"] if worker else 2224
|
||||
actions.append(f"Analyzing pre-flight health state for {agent_name} (port {port})")
|
||||
|
||||
# 1. Ensure Gitea Collaborator & Partition Table
|
||||
token = ""
|
||||
if PARTITION_TABLE_PATH.exists():
|
||||
try:
|
||||
with open(PARTITION_TABLE_PATH) as f:
|
||||
pt = json.load(f)
|
||||
contributor = pt.get("contributors", {}).get(agent_name)
|
||||
if contributor:
|
||||
token = contributor.get("token", "")
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
if not token:
|
||||
token = hashlib.sha256(f"{agent_name}-gitea-token".encode()).hexdigest()[:40]
|
||||
actions.append(f"Generated partition token for {agent_name}")
|
||||
|
||||
collab_res = gitea_api_request(f"/repos/super/box/collaborators/{agent_name}", method="PUT", data={"permission": "write"})
|
||||
actions.append(f"Ensured Gitea collaborator write access for {agent_name}")
|
||||
|
||||
# 2. Container Workspace Injection if SSH dialable
|
||||
if port in (2224, 2228):
|
||||
try:
|
||||
cmd = f"ssh -o ConnectTimeout=3 -o BatchMode=yes -o StrictHostKeyChecking=no super@100.81.31.9 'ssh -o StrictHostKeyChecking=no -i /home/super/.ssh/fleet -p {port} muse@localhost \"git config --global credential.helper store && echo \\\"https://{agent_name}:{token}@tea.muse-dev.online\\\" > ~/.git-credentials && chmod 600 ~/.git-credentials\"' 2>/dev/null"
|
||||
if os.system(cmd) == 0:
|
||||
actions.append(f"Directly injected Git credentials into {agent_name} container")
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# 3. Check and heal Hatch / Reverse Tunnel
|
||||
tunnel_ports = check_tunnel_ports()
|
||||
if tunnel_ports.get(port) == "UP":
|
||||
actions.append(f"Hatch reverse tunnel verified UP on port {port}")
|
||||
else:
|
||||
chat_script = REPO_ROOT / "bin" / "muse-chat-api.py"
|
||||
if chat_script.exists():
|
||||
heal_msg = f"[HEAL NUDGE] Reverse tunnel on port {port} is DOWN. Please run 'chmod 600 ~/.ssh/authorized_keys' and restart tunnel with '~/workspace/bin/gcp-tunnel-up.sh &' (or 'cloud-uptime/recover-after-rebuild.sh'). Git clone URL: https://{agent_name}:{token}@tea.muse-dev.online/super/box.git"
|
||||
os.system(f"python3 {chat_script} --account {agent_name} send '{heal_msg}' >/dev/null 2>&1")
|
||||
actions.append(f"Dispatched tunnel restart & git clone command to {agent_name} chat")
|
||||
|
||||
time.sleep(1.0)
|
||||
recheck_ports = check_tunnel_ports()
|
||||
if recheck_ports.get(port) == "UP":
|
||||
actions.append(f"Reverse tunnel on port {port} came online during healing!")
|
||||
else:
|
||||
unresolved.append(f"Reverse tunnel on port {port} is still DOWN (waiting for agent container execution)")
|
||||
|
||||
final_preflight = check_agent_preflight(agent_name)
|
||||
healed = final_preflight["ready"]
|
||||
if not healed and not unresolved:
|
||||
unresolved.extend(final_preflight["reasons"])
|
||||
|
||||
return {
|
||||
"agent": agent_name,
|
||||
"healed": healed,
|
||||
"actions": actions,
|
||||
"unresolved": unresolved
|
||||
}
|
||||
|
||||
def cmd_heal(args):
|
||||
agent = args.agent
|
||||
print(c_bold(f"\n=== BOX WORK: HEALING AGENT '{agent}' ===\n"))
|
||||
res = heal_agent(agent)
|
||||
print(c_bold("Actions taken:"))
|
||||
for a in res["actions"]:
|
||||
print(f" {c_green('✓')} {a}")
|
||||
print()
|
||||
if res["healed"]:
|
||||
print(c_green(f"🎉 Agent '{agent}' successfully healed and ready for assignments!\n"))
|
||||
else:
|
||||
print(c_yellow(f"⚠️ Agent '{agent}' partially healed with open issues:"))
|
||||
for u in res["unresolved"]:
|
||||
print(f" • {u}")
|
||||
print()
|
||||
|
||||
def cmd_status(args):
|
||||
api_base, _ = get_gitea_config()
|
||||
print(c_bold(f"\n=== BOX WORK: FLEET & BUILD PIPELINE ({api_base}) ===\n"))
|
||||
|
||||
@@ -550,18 +643,26 @@ def cmd_start(args):
|
||||
agent = args.agent
|
||||
body = args.goal or f"Work task for {agent}: {title}"
|
||||
|
||||
# 0. Pre-flight health gate: Hatch, Restore, Git Config
|
||||
# 0. Pre-flight health gate: Hatch, Restore, Git Config with Auto-Heal
|
||||
preflight = check_agent_preflight(agent)
|
||||
if not preflight["ready"] and not getattr(args, "force", False):
|
||||
print(c_red(f"\n[BLOCKED] Agent '{agent}' failed pre-flight health verification:"))
|
||||
print(c_yellow(f"\n[PRE-FLIGHT FAILED] Agent '{agent}' requires healing before assignment."))
|
||||
print(f" • Hatch: {badge_status(preflight['hatch']['status'])} - {preflight['hatch']['details']}")
|
||||
print(f" • Restore: {badge_status(preflight['restore']['status'])} - {preflight['restore']['details']}")
|
||||
print(f" • Git: {badge_status(preflight['git']['status'])} - {preflight['git']['details']}")
|
||||
print(c_yellow("\nBlocking reasons:"))
|
||||
for r in preflight["reasons"]:
|
||||
print(f" - {r}")
|
||||
print(c_dim(f"\nTo inspect full health: box work check {agent}\nTo bypass pre-flight: box work start '{title}' --to {agent} --force\n"))
|
||||
sys.exit(1)
|
||||
print(c_bold("\nAttempting automated remediation (auto-heal)..."))
|
||||
heal_res = heal_agent(agent)
|
||||
for a in heal_res["actions"]:
|
||||
print(f" {c_green('✓')} {a}")
|
||||
|
||||
if heal_res["healed"]:
|
||||
print(c_green(f"\n🎉 Successfully healed {agent}! Proceeding with ticket dispatch..."))
|
||||
else:
|
||||
print(c_red(f"\n[BLOCKED] Auto-heal could not resolve all issues for {agent}:"))
|
||||
for issue in heal_res["unresolved"]:
|
||||
print(f" • {issue}")
|
||||
print(c_dim(f"\nTo inspect: box work check {agent}\nTo bypass: box work start '{title}' --to {agent} --force\n"))
|
||||
sys.exit(1)
|
||||
elif not preflight["ready"] and getattr(args, "force", False):
|
||||
print(c_yellow(f"[WARNING] Overriding failed pre-flight checks on {agent} (--force specified).\n"))
|
||||
else:
|
||||
@@ -670,9 +771,7 @@ def cmd_merge(args):
|
||||
|
||||
def cmd_chats(args):
|
||||
agent = getattr(args, "agent", None)
|
||||
events = get_recent_chat_events(limit=args.limit)
|
||||
if agent:
|
||||
events = [e for e in events if e.get("agent") == agent]
|
||||
events = get_recent_chat_events(limit=args.limit, agent=agent)
|
||||
print(c_bold(f"\n=== CHAT FEED ({agent or 'ALL AGENTS'}) ===\n"))
|
||||
for ev in events:
|
||||
ag = ev.get("agent", "agent")
|
||||
|
||||
+19
-5
@@ -7,7 +7,7 @@ Usage:
|
||||
Requires: websocket-client (pip install --break-system-packages websocket-client)
|
||||
CDP relay must be up: http://10.201.87.2:9410/json/list
|
||||
"""
|
||||
import json, sys, time, urllib.request
|
||||
import json, sys, time, urllib.request, subprocess, os
|
||||
import websocket
|
||||
|
||||
CDP_URL = "http://10.201.87.2:9410/json/list"
|
||||
@@ -20,10 +20,24 @@ def get_page():
|
||||
raise RuntimeError("no muse.ai page found")
|
||||
return pages[0]
|
||||
|
||||
def connect():
|
||||
page = get_page()
|
||||
ws = websocket.create_connection(page['webSocketDebuggerUrl'], timeout=10)
|
||||
return ws
|
||||
def revive_browser():
|
||||
script = os.path.expanduser("~/Projects/NetVM/bin/netvm-chrome.sh")
|
||||
print("CDP disconnect detected. Reviving Chromium smoke profile...", file=sys.stderr)
|
||||
subprocess.Popen([script, "--headless", "smoke", "https://muse.ai"], stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL)
|
||||
time.sleep(4)
|
||||
|
||||
def connect(retries=2):
|
||||
for attempt in range(retries + 1):
|
||||
try:
|
||||
page = get_page()
|
||||
ws = websocket.create_connection(page['webSocketDebuggerUrl'], timeout=10)
|
||||
return ws
|
||||
except Exception as e:
|
||||
if attempt < retries:
|
||||
revive_browser()
|
||||
else:
|
||||
raise RuntimeError(f"Failed to connect to CDP after {retries} retries: {e}")
|
||||
|
||||
|
||||
def ev(ws, expr, await_promise=False):
|
||||
ws.send(json.dumps({"id":1,"method":"Runtime.evaluate",
|
||||
|
||||
Reference in New Issue
Block a user