#!/usr/bin/env python3 """box-readback-loopback.py — Synchronous Readback Gate & Cognitive-Aware True Loopback Engine. Architecture: 1. Readback Gate: - Synchronously awaits agent confirmation (up to 45s) after task dispatch. - Validates via Hybrid Tag ([READBACK] Ticket #...) + Semantic Fallback. - Marks ticket 'in-progress' upon confirmation; marks 'blocked' and releases agent on timeout. 2. Cognitive-Aware True Loopbacks: - Zero Token Burn: 100% silent while git commits or PRs are progressing. - Cognitive Guard: Defer loopbacks while agent is THINKING / GENERATING. - Dual-Layer Escalation: * 15m inactive: Tier 1 non-intrusive Gitea ticket comment (@agent inquiry). * 45m inactive: Tier 2 direct chat DM escalation (muse-cli-node send). * 90m inactive: Tier 3 failure escalation (mark 'blocked', alert #lobby, release agent). """ import json import os import re import socket import subprocess import sys import time import urllib.request from datetime import datetime, timezone from pathlib import Path # Local imports try: import agent_cognitive_probe as acp except ImportError: acp = None REMOTE_HOST = "100.123.153.75" # bl control node DEFAULT_GITEA_URL = "https://tea.muse-dev.online" LOOPBACK_STATE_FILE = Path("/tmp/box-loopback-state.json") def is_running_on_bl(): try: hn = socket.gethostname().lower() if "bl" in hn: return True except Exception: pass return os.path.exists("/var/run/netns/warp-muse") or os.path.exists("/run/netns/warp-muse") def get_gitea_token(): token = os.environ.get("GITEA_TOKEN", "3c26744525bceaf385aa09737f7e41af613627b6") return token def gitea_api(endpoint: str, method: str = "GET", data: dict = None): token = get_gitea_token() url = f"{DEFAULT_GITEA_URL}/api/v1{endpoint}" headers = { "Authorization": f"token {token}", "Content-Type": "application/json", "User-Agent": "Box-Work-CLI/1.0", } payload = json.dumps(data).encode("utf-8") if data else None req = urllib.request.Request(url, data=payload, headers=headers, method=method) try: with urllib.request.urlopen(req, timeout=10) as r: if r.status in (200, 201): return json.loads(r.read().decode()) return {"status": r.status} except urllib.error.HTTPError as e: try: return json.loads(e.read().decode()) except Exception: return {"error": str(e), "code": e.code} except Exception as e: return {"error": str(e)} # ---------------------------------------------------------------------- # CHAT COMMUNICATION HELPERS # ---------------------------------------------------------------------- def send_agent_chat(agent: str, message: str) -> bool: """Delivers a message directly into the agent's web chat session.""" if is_running_on_bl(): cmd = ["/home/super/Projects/NetVM/bin/muse-cli-node", agent, "send", message] else: cmd = ["ssh", "-q", f"super@{REMOTE_HOST}", f"/home/super/Projects/NetVM/bin/muse-cli-node {agent} send {subprocess.list2cmdline([message])}"] try: res = subprocess.run(cmd, capture_output=True, text=True, timeout=15) return res.returncode == 0 except Exception: return False def get_agent_history(agent: str, limit: int = 5) -> list: """Retrieves recent chat messages from the agent's active session.""" if is_running_on_bl(): cmd = ["/home/super/Projects/NetVM/bin/muse-cli-node", agent, "history", "--limit", str(limit)] else: cmd = ["ssh", "-q", f"super@{REMOTE_HOST}", f"/home/super/Projects/NetVM/bin/muse-cli-node {agent} history --limit {limit}"] try: res = subprocess.run(cmd, capture_output=True, text=True, timeout=15) if res.returncode == 0 and res.stdout.strip(): data = json.loads(res.stdout.strip()) if isinstance(data, list): return data except Exception: pass return [] def get_latest_chat_seq(agent: str) -> int: """Finds the maximum sequence number in the agent's chat history.""" history = get_agent_history(agent, limit=3) seqs = [m.get("seq", 0) for m in history if isinstance(m, dict) and "seq" in m] return max(seqs) if seqs else 0 # ---------------------------------------------------------------------- # READBACK GATE # ---------------------------------------------------------------------- def validate_readback(text: str, issue_num: int, agent: str) -> tuple[bool, str]: """Validates an agent readback using Hybrid Tag + Semantic Fallback. Returns (is_valid, excerpt). """ if not text: return False, "" clean_text = text.strip() issue_pattern = rf"#?{issue_num}\b" # 1. Strict Tag Match: [READBACK] Ticket # ... if re.search(r"\[READBACK\]", clean_text, re.IGNORECASE) and re.search(issue_pattern, clean_text): snippet = clean_text[:200].replace("\n", " ") return True, snippet # 2. Semantic Fallback: Mentions ticket number AND branch/accepted status has_issue = bool(re.search(issue_pattern, clean_text)) has_branch_or_ack = bool(re.search( rf"(dev/{agent}/|branch|accepted|working on|confirm|start(ed|ing)|received)", clean_text, re.IGNORECASE )) if has_issue and has_branch_or_ack: snippet = clean_text[:200].replace("\n", " ") return True, snippet return False, "" def wait_for_readback(agent: str, issue_num: int, initial_seq: int, timeout_s: int = 45, poll_s: float = 3.0) -> dict: """Synchronously polls for agent readback within timeout_s.""" start_time = time.time() deadline = start_time + timeout_s while time.time() < deadline: elapsed = int(time.time() - start_time) print(f"\r ⏳ Awaiting Readback from @{agent} ({elapsed}s / {timeout_s}s)...", end="", flush=True) history = get_agent_history(agent, limit=4) for msg in history: seq = msg.get("seq", 0) role = msg.get("role", "") text = msg.get("text", "") # Only check new assistant messages if seq > initial_seq and role == "assistant": valid, excerpt = validate_readback(text, issue_num, agent) if valid: print() return { "success": True, "snippet": excerpt, "elapsed": elapsed, "seq": seq, } time.sleep(poll_s) print() return { "success": False, "timeout": True, "elapsed": timeout_s, } def handle_readback_success(agent: str, issue_num: int, snippet: str): """Marks ticket in-progress and records confirmation on Gitea.""" # Label ticket in-progress gitea_api(f"/repos/super/box/issues/{issue_num}/labels", method="POST", data={"labels": ["in-progress"]}) # Post confirmation comment comment_body = f"🤖 **Readback Confirmed** by @{agent}:\n> {snippet}" gitea_api(f"/repos/super/box/issues/{issue_num}/comments", method="POST", data={"body": comment_body}) def handle_readback_timeout(agent: str, issue_num: int, title: str): """Labels ticket blocked and unassigns agent so they return to IDLE.""" # Label ticket blocked gitea_api(f"/repos/super/box/issues/{issue_num}/labels", method="POST", data={"labels": ["blocked"]}) # Post explanation comment comment_body = ( f"⚠️ **Readback Timeout**: Agent @{agent} did not confirm ticket #{issue_num} " f"within 45 seconds of dispatch. Releasing assignment to prevent deadlocks." ) gitea_api(f"/repos/super/box/issues/{issue_num}/comments", method="POST", data={"body": comment_body}) # Unassign agent gitea_api(f"/repos/super/box/issues/{issue_num}", method="PATCH", data={"assignees": []}) # ---------------------------------------------------------------------- # COGNITIVE-AWARE TRUE LOOPBACK ENGINE # ---------------------------------------------------------------------- def load_loopback_state() -> dict: if LOOPBACK_STATE_FILE.exists(): try: with open(LOOPBACK_STATE_FILE) as f: return json.load(f) except Exception: pass return {} def save_loopback_state(state: dict): try: with open(LOOPBACK_STATE_FILE, "w") as f: json.dump(state, f, indent=2) except Exception: pass def get_ticket_git_activity(agent: str, issue_num: int) -> datetime | None: """Checks the latest commit timestamp on the agent's branch dev//-*.""" # Check Gitea branches for dev//-* branches = gitea_api("/repos/super/box/branches") if isinstance(branches, list): target_prefix = f"dev/{agent}/{issue_num}" for b in branches: name = b.get("name", "") if target_prefix in name: commit = b.get("commit", {}) ts_str = commit.get("timestamp") if ts_str: try: return datetime.fromisoformat(ts_str.replace("Z", "+00:00")) except Exception: pass return None def get_ticket_last_activity(issue: dict, agent: str) -> tuple[datetime, str]: """Finds the most recent activity timestamp (git commit, comment, or issue creation).""" issue_num = issue["number"] latest_dt = datetime.fromisoformat(issue["created_at"].replace("Z", "+00:00")) source = "issue_created" # Check comments comments = gitea_api(f"/repos/super/box/issues/{issue_num}/comments") if isinstance(comments, list): for c in comments: c_dt = datetime.fromisoformat(c["created_at"].replace("Z", "+00:00")) if c_dt > latest_dt: latest_dt = c_dt source = "gitea_comment" # Check git branch commit git_dt = get_ticket_git_activity(agent, issue_num) if git_dt and git_dt > latest_dt: latest_dt = git_dt source = "git_commit" return latest_dt, source def run_loopback_sweep(dry_run: bool = False, verbose: bool = True) -> list: """Executes a single sweep of all open assigned tickets according to the 3-tier escalation model.""" now = datetime.now(timezone.utc) state = load_loopback_state() actions_taken = [] issues = gitea_api("/repos/super/box/issues?state=open") if not isinstance(issues, list): if verbose: print("Failed to fetch open issues from Gitea.") return [] assigned_issues = [i for i in issues if i.get("assignee")] if verbose: print(f"\n=== LOOPBACK SWEEP: {len(assigned_issues)} ACTIVE ASSIGNED TICKETS ({now.strftime('%H:%M:%SZ')}) ===") for iss in assigned_issues: issue_num = iss["number"] title = iss.get("title", "") agent = iss["assignee"]["username"] key = str(issue_num) ticket_state = state.get(key, {}) last_dt, source = get_ticket_last_activity(iss, agent) inactive_s = (now - last_dt).total_seconds() inactive_m = int(inactive_s // 60) # Check cognitive state cog = acp.get_passive_cognitive_state(agent) if acp else {"status": "IDLE", "cognitive_lock": False} cog_status = cog.get("status", "IDLE") is_thinking = cog_status == "THINKING" or cog.get("cognitive_lock") if verbose: print(f"Ticket #{issue_num} (@{agent}): {inactive_m}m inactive (source: {source}) | Cognitive: {cog_status}") # Tier 0: Inactive < 15m or active git commits -> Complete silence if inactive_m < 15 or source == "git_commit": if verbose: print(" 👉 Status: Active or within silent grace period (<15m). No action.") continue # Check cognitive guard: Defer if agent is thinking/generating if is_thinking and cog_status != "INPUT_WAIT": if verbose: print(f" 🧠 Cognitive Guard: Deferring loopback — @{agent} is currently {cog_status}.") continue # Tier 1: 15m <= Inactivity < 45m -> Non-intrusive Gitea ticket comment if 15 <= inactive_m < 45: if ticket_state.get("tier1_sent"): if verbose: print(" 👉 Tier 1 comment already dispatched. Waiting for 45m threshold.") continue msg = ( f"🤖 @{agent} **Loopback Tier 1 Check** ({inactive_m}m elapsed):\n" f"No git commits recorded on feature branch for Ticket #{issue_num}. " f"Are you progressing or blocked? Reply with status or push a commit." ) action_desc = f"Tier 1: Posted Gitea comment to #{issue_num} (@{agent})" actions_taken.append(action_desc) if not dry_run: gitea_api(f"/repos/super/box/issues/{issue_num}/comments", method="POST", data={"body": msg}) ticket_state["tier1_sent"] = now.isoformat() state[key] = ticket_state save_loopback_state(state) if verbose: print(f" ✓ {action_desc}") # Tier 2: 45m <= Inactivity < 90m -> Direct Chat DM Escalation elif 45 <= inactive_m < 90: if ticket_state.get("tier2_sent"): if verbose: print(" 👉 Tier 2 chat DM already dispatched. Waiting for 90m threshold.") continue chat_msg = ( f"[LOOPBACK ALERT] Ticket #{issue_num} ('{title}'): " f"{inactive_m} minutes inactive with no git commits. " f"Please confirm if blocked on tool execution, terminal approvals, or environment." ) action_desc = f"Tier 2: Escalated to chat DM for @{agent} on #{issue_num}" actions_taken.append(action_desc) if not dry_run: send_agent_chat(agent, chat_msg) gitea_api(f"/repos/super/box/issues/{issue_num}/comments", method="POST", data={ "body": f"📣 **Loopback Tier 2 Escalation**: Inactivity reached {inactive_m}m. Sent direct chat DM to @{agent}." }) ticket_state["tier2_sent"] = now.isoformat() state[key] = ticket_state save_loopback_state(state) if verbose: print(f" ✓ {action_desc}") # Tier 3: Inactivity >= 90m (or unhandled INPUT_WAIT > 15m) -> Fail-closed escalation elif inactive_m >= 90 or (cog_status == "INPUT_WAIT" and inactive_m >= 15): action_desc = f"Tier 3: Ticket #{issue_num} marked BLOCKED; released @{agent} assignment" actions_taken.append(action_desc) if not dry_run: # Label blocked gitea_api(f"/repos/super/box/issues/{issue_num}/labels", method="POST", data={"labels": ["blocked"]}) # Post failure comment gitea_api(f"/repos/super/box/issues/{issue_num}/comments", method="POST", data={ "body": ( f"🚨 **Loopback Tier 3 Escalation**: Inactivity reached {inactive_m}m with zero git progress. " f"Ticket marked `blocked` and unassigned from @{agent} for operator intervention." ) }) # Unassign agent gitea_api(f"/repos/super/box/issues/{issue_num}", method="PATCH", data={"assignees": []}) ticket_state["tier3_sent"] = now.isoformat() state[key] = ticket_state save_loopback_state(state) if verbose: print(f" 🚨 {action_desc}") if verbose: print() return actions_taken def main(): import argparse parser = argparse.ArgumentParser(description="Synchronous Readback & Cognitive True Loopback Engine") sub = parser.add_subparsers(dest="cmd") p_sweep = sub.add_parser("sweep", help="Run a loopback sweep across open tickets") p_sweep.add_argument("--dry-run", action="store_true", help="Evaluate conditions without sending messages") p_sweep.add_argument("--json", action="store_true", help="Output actions as JSON") args = parser.parse_args() if not args.cmd or args.cmd == "sweep": dry_run = getattr(args, "dry_run", False) as_json = getattr(args, "json", False) actions = run_loopback_sweep(dry_run=dry_run, verbose=not as_json) if as_json: print(json.dumps({"ok": True, "actions": actions}, indent=2)) if __name__ == "__main__": main()