diff --git a/bin/box-work.py b/bin/box-work.py index 58bb6f5..46ec5c1 100755 --- a/bin/box-work.py +++ b/bin/box-work.py @@ -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") diff --git a/bin/muse-chat-api.py b/bin/muse-chat-api.py index d620b47..9ba94a3 100755 --- a/bin/muse-chat-api.py +++ b/bin/muse-chat-api.py @@ -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",