Files
box/bin/box-readback-loopback.py
T

425 lines
16 KiB
Python
Raw Normal View History

#!/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 #<num> ...
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/<agent>/<issue>-*."""
# Check Gitea branches for dev/<agent>/<issue_num>-*
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()