feat(work): implement synchronous readback gate and cognitive true loopback engine
This commit is contained in:
Executable
+424
@@ -0,0 +1,424 @@
|
|||||||
|
#!/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()
|
||||||
+69
-12
@@ -31,6 +31,11 @@ try:
|
|||||||
except ImportError:
|
except ImportError:
|
||||||
acp = None
|
acp = None
|
||||||
|
|
||||||
|
try:
|
||||||
|
import box_readback_loopback as brl
|
||||||
|
except ImportError:
|
||||||
|
brl = None
|
||||||
|
|
||||||
# Color helpers
|
# Color helpers
|
||||||
USE_COLOR = sys.stdout.isatty() or os.environ.get("CLICOLOR_FORCE") == "1"
|
USE_COLOR = sys.stdout.isatty() or os.environ.get("CLICOLOR_FORCE") == "1"
|
||||||
|
|
||||||
@@ -741,18 +746,47 @@ def cmd_start(args):
|
|||||||
pass
|
pass
|
||||||
s.close()
|
s.close()
|
||||||
|
|
||||||
# 4. Notify agent via muse-chat-api if available
|
# 4. Notify agent via direct chat with Readback requirement
|
||||||
chat_script = REPO_ROOT / "bin" / "muse-chat-api.py"
|
prompt_msg = (
|
||||||
if chat_script.exists():
|
f"New build ticket #{issue_num} assigned to you: {title}.\n"
|
||||||
msg = f"New build ticket #{issue_num} assigned to you: {title}. Clone/pull ~/workspace/box, checkout dev/{agent}/{issue_num}-work, commit citing 'Fixes #{issue_num}', and push."
|
f"Goal: {body}\n\n"
|
||||||
try:
|
f"Please confirm with Readback:\n"
|
||||||
cmd = f"python3 {chat_script} --account {agent} send '{msg}'"
|
f"[READBACK] Ticket #{issue_num} | Branch: dev/{agent}/{issue_num}-work | Plan: <short summary>"
|
||||||
os.system(f"{cmd} >/dev/null 2>&1")
|
)
|
||||||
print(c_green(f"✓ Delivered briefing to {agent} chat session"))
|
|
||||||
except Exception:
|
|
||||||
pass
|
|
||||||
|
|
||||||
print(c_bold(f"\nWork ticket #{issue_num} is active and assigned to {agent}.\n"))
|
initial_seq = brl.get_latest_chat_seq(agent) if brl else 0
|
||||||
|
delivered = brl.send_agent_chat(agent, prompt_msg) if brl else False
|
||||||
|
if delivered:
|
||||||
|
print(c_green(f"✓ Delivered briefing to @{agent} chat session"))
|
||||||
|
else:
|
||||||
|
chat_script = REPO_ROOT / "bin" / "muse-chat-api.py"
|
||||||
|
if chat_script.exists():
|
||||||
|
os.system(f"python3 {chat_script} --account {agent} send '{prompt_msg}' >/dev/null 2>&1")
|
||||||
|
print(c_green(f"✓ Delivered briefing to @{agent} chat session (via muse-chat-api)"))
|
||||||
|
else:
|
||||||
|
print(c_yellow(f"⚠️ Direct chat delivery unconfirmed for @{agent}"))
|
||||||
|
|
||||||
|
# 5. Synchronous Readback Gate (Wait up to 45s)
|
||||||
|
if getattr(args, "no_readback", False):
|
||||||
|
print(c_yellow(f"[OVERRIDE] Skipping readback gate (--no-readback specified)."))
|
||||||
|
print(c_bold(f"\nWork ticket #{issue_num} is active and assigned to {agent}.\n"))
|
||||||
|
elif brl:
|
||||||
|
print(c_bold(f"\nEngaging Synchronous Readback Gate (timeout: 45s)..."))
|
||||||
|
rb_res = brl.wait_for_readback(agent, issue_num, initial_seq, timeout_s=45)
|
||||||
|
if rb_res["success"]:
|
||||||
|
brl.handle_readback_success(agent, issue_num, rb_res["snippet"])
|
||||||
|
print(c_green(f"\n🎉 Readback Confirmed from @{agent} in {rb_res['elapsed']}s!"))
|
||||||
|
print(c_dim(f" Excerpt: {rb_res['snippet'][:120]}..."))
|
||||||
|
print(c_bold(f"\nWork ticket #{issue_num} is confirmed in-progress and assigned to {agent}.\n"))
|
||||||
|
else:
|
||||||
|
brl.handle_readback_timeout(agent, issue_num, title)
|
||||||
|
print(c_red(f"\n❌ [READBACK TIMEOUT] Agent @{agent} failed to confirm Ticket #{issue_num} within 45s."))
|
||||||
|
print(c_yellow(f" • Ticket marked 'blocked' in Gitea"))
|
||||||
|
print(c_yellow(f" • Assignment released so @{agent} returns to IDLE"))
|
||||||
|
print(c_dim(f"\nTo retry: box work start '{title}' --to {agent}\nTo bypass: box work start '{title}' --to {agent} --no-readback\n"))
|
||||||
|
sys.exit(1)
|
||||||
|
else:
|
||||||
|
print(c_bold(f"\nWork ticket #{issue_num} is active and assigned to {agent}.\n"))
|
||||||
|
|
||||||
def cmd_assign(args):
|
def cmd_assign(args):
|
||||||
issue_num = args.issue
|
issue_num = args.issue
|
||||||
@@ -863,18 +897,34 @@ def cmd_cognitive(args):
|
|||||||
sys.exit(1)
|
sys.exit(1)
|
||||||
acp.cmd_status(args)
|
acp.cmd_status(args)
|
||||||
|
|
||||||
|
def cmd_loopback(args):
|
||||||
|
if not brl:
|
||||||
|
print(c_red("Error: box_readback_loopback module not found."))
|
||||||
|
sys.exit(1)
|
||||||
|
dry_run = getattr(args, "dry_run", False)
|
||||||
|
as_json = getattr(args, "json", False)
|
||||||
|
actions = brl.run_loopback_sweep(dry_run=dry_run, verbose=not as_json)
|
||||||
|
if as_json:
|
||||||
|
print(json.dumps({"ok": True, "actions": actions}, indent=2))
|
||||||
|
|
||||||
WORK_COMMAND_EXAMPLES = {
|
WORK_COMMAND_EXAMPLES = {
|
||||||
"box work": [
|
"box work": [
|
||||||
"box work # View fleet workspace dashboard & signals",
|
"box work # View fleet workspace dashboard & signals",
|
||||||
|
"box work loopback [--dry-run] # Sweep open tickets with 3-tier loopback escalation",
|
||||||
"box work cognitive [agent...] # Live zero-click cognitive sensor probe across fleet",
|
"box work cognitive [agent...] # Live zero-click cognitive sensor probe across fleet",
|
||||||
"box work menu <agent> [tab] # Inspect agent profile menu (tasks, timers, approvals)",
|
"box work menu <agent> [tab] # Inspect agent profile menu (tasks, timers, approvals)",
|
||||||
"box work check [agent] # Audit pre-flight health gates",
|
"box work check [agent] # Audit pre-flight health gates",
|
||||||
"box work heal <agent> # Automated remediation & chat nudge",
|
"box work heal <agent> # Automated remediation & chat nudge",
|
||||||
"box work start \"<title>\" --to <agent> # Start & dispatch new build ticket",
|
"box work start \"<title>\" --to <agent> # Start & dispatch new build ticket with readback gate",
|
||||||
"box work assign <issue#> --to <agent> # Assign existing ticket",
|
"box work assign <issue#> --to <agent> # Assign existing ticket",
|
||||||
"box work merge <pr#> # Verify tests and merge PR to master",
|
"box work merge <pr#> # Verify tests and merge PR to master",
|
||||||
"box work chats --agent <name> # View live multi-agent chat feed",
|
"box work chats --agent <name> # View live multi-agent chat feed",
|
||||||
],
|
],
|
||||||
|
"box work loopback": [
|
||||||
|
"box work loopback # Sweep active tickets and escalate inactivity",
|
||||||
|
"box work loopback --dry-run # Dry-run inspect without sending pings or comments",
|
||||||
|
"box work loopback --json # Machine-readable output for systemd cron",
|
||||||
|
],
|
||||||
"box work menu": [
|
"box work menu": [
|
||||||
"box work menu muse upcoming # Inspect timers & recurring cron loops",
|
"box work menu muse upcoming # Inspect timers & recurring cron loops",
|
||||||
"box work menu 646 activity # Inspect recent tasks & active processes",
|
"box work menu 646 activity # Inspect recent tasks & active processes",
|
||||||
@@ -1000,6 +1050,7 @@ def main():
|
|||||||
p_start.add_argument("--goal", help="Optional detailed goal description")
|
p_start.add_argument("--goal", help="Optional detailed goal description")
|
||||||
p_start.add_argument("--force", action="store_true", help="Bypass pre-flight health gate")
|
p_start.add_argument("--force", action="store_true", help="Bypass pre-flight health gate")
|
||||||
p_start.add_argument("--no-heal", action="store_true", help="Fail immediately without attempting auto-heal if pre-flight checks fail")
|
p_start.add_argument("--no-heal", action="store_true", help="Fail immediately without attempting auto-heal if pre-flight checks fail")
|
||||||
|
p_start.add_argument("--no-readback", action="store_true", help="Bypass synchronous 45s readback gate")
|
||||||
|
|
||||||
p_assign = sub.add_parser("assign", help="Assign existing ticket to an agent")
|
p_assign = sub.add_parser("assign", help="Assign existing ticket to an agent")
|
||||||
p_assign.add_argument("issue", type=int, help="Issue number (e.g. 215)")
|
p_assign.add_argument("issue", type=int, help="Issue number (e.g. 215)")
|
||||||
@@ -1026,6 +1077,10 @@ def main():
|
|||||||
p_cog.add_argument("agents", nargs="*", help="Optional agent usernames")
|
p_cog.add_argument("agents", nargs="*", help="Optional agent usernames")
|
||||||
p_cog.add_argument("--json", action="store_true", help="Output JSON")
|
p_cog.add_argument("--json", action="store_true", help="Output JSON")
|
||||||
|
|
||||||
|
p_loop = sub.add_parser("loopback", help="Sweep active tickets with 3-tier loopback escalation")
|
||||||
|
p_loop.add_argument("--dry-run", action="store_true", help="Inspect without modifying tickets or sending DMs")
|
||||||
|
p_loop.add_argument("--json", action="store_true", help="Output actions as JSON")
|
||||||
|
|
||||||
args = parser.parse_args()
|
args = parser.parse_args()
|
||||||
action = args.work_action
|
action = args.work_action
|
||||||
|
|
||||||
@@ -1047,6 +1102,8 @@ def main():
|
|||||||
cmd_menu(args)
|
cmd_menu(args)
|
||||||
elif action == "cognitive":
|
elif action == "cognitive":
|
||||||
cmd_cognitive(args)
|
cmd_cognitive(args)
|
||||||
|
elif action == "loopback":
|
||||||
|
cmd_loopback(args)
|
||||||
else:
|
else:
|
||||||
parser.print_help()
|
parser.print_help()
|
||||||
|
|
||||||
|
|||||||
Symlink
+1
@@ -0,0 +1 @@
|
|||||||
|
box-readback-loopback.py
|
||||||
@@ -1493,6 +1493,8 @@ def cmd_work(args):
|
|||||||
box_work.cmd_menu(args)
|
box_work.cmd_menu(args)
|
||||||
elif action == "cognitive":
|
elif action == "cognitive":
|
||||||
box_work.cmd_cognitive(args)
|
box_work.cmd_cognitive(args)
|
||||||
|
elif action == "loopback":
|
||||||
|
box_work.cmd_loopback(args)
|
||||||
else:
|
else:
|
||||||
box_work.cmd_status(args)
|
box_work.cmd_status(args)
|
||||||
|
|
||||||
@@ -6998,6 +7000,7 @@ def build_parser():
|
|||||||
p_w_start.add_argument("--goal", help="Optional detailed goal description")
|
p_w_start.add_argument("--goal", help="Optional detailed goal description")
|
||||||
p_w_start.add_argument("--force", action="store_true", help="Bypass pre-flight health gate")
|
p_w_start.add_argument("--force", action="store_true", help="Bypass pre-flight health gate")
|
||||||
p_w_start.add_argument("--no-heal", action="store_true", help="Fail immediately without attempting auto-heal if pre-flight checks fail")
|
p_w_start.add_argument("--no-heal", action="store_true", help="Fail immediately without attempting auto-heal if pre-flight checks fail")
|
||||||
|
p_w_start.add_argument("--no-readback", action="store_true", help="Bypass synchronous 45s readback gate")
|
||||||
p_w_assign = work_sub.add_parser("assign", parents=[common], help="Assign existing ticket to an agent")
|
p_w_assign = work_sub.add_parser("assign", parents=[common], help="Assign existing ticket to an agent")
|
||||||
p_w_assign.add_argument("issue", type=int, help="Issue number (e.g. 215)")
|
p_w_assign.add_argument("issue", type=int, help="Issue number (e.g. 215)")
|
||||||
p_w_assign.add_argument("--to", dest="agent", required=True, help="Agent username")
|
p_w_assign.add_argument("--to", dest="agent", required=True, help="Agent username")
|
||||||
@@ -7013,6 +7016,8 @@ def build_parser():
|
|||||||
p_w_menu.add_argument("tab", nargs="?", default="activity", choices=["activity", "upcoming", "approvals", "identity", "all"], help="Menu tab to view")
|
p_w_menu.add_argument("tab", nargs="?", default="activity", choices=["activity", "upcoming", "approvals", "identity", "all"], help="Menu tab to view")
|
||||||
p_w_cog = work_sub.add_parser("cognitive", parents=[common], help="Probe real-time cognitive sensor (thinking, generating, sidechats)")
|
p_w_cog = work_sub.add_parser("cognitive", parents=[common], help="Probe real-time cognitive sensor (thinking, generating, sidechats)")
|
||||||
p_w_cog.add_argument("agents", nargs="*", help="Optional agent usernames")
|
p_w_cog.add_argument("agents", nargs="*", help="Optional agent usernames")
|
||||||
|
p_w_loop = work_sub.add_parser("loopback", parents=[common], help="Sweep active tickets with 3-tier loopback escalation")
|
||||||
|
p_w_loop.add_argument("--dry-run", action="store_true", help="Inspect without modifying tickets or sending DMs")
|
||||||
|
|
||||||
p_tasks = subparsers.add_parser("tasks", parents=[common], help="Agent work queue: pending/claimed/done files (distinct from scheduled jobs)")
|
p_tasks = subparsers.add_parser("tasks", parents=[common], help="Agent work queue: pending/claimed/done files (distinct from scheduled jobs)")
|
||||||
p_tasks.add_argument("--dir", default=None, help="Task queue dir (default: fleet/tasks)")
|
p_tasks.add_argument("--dir", default=None, help="Task queue dir (default: fleet/tasks)")
|
||||||
|
|||||||
@@ -0,0 +1,10 @@
|
|||||||
|
[Unit]
|
||||||
|
Description=Box Work Cognitive-Aware True Loopback Sweeper
|
||||||
|
After=network.target
|
||||||
|
|
||||||
|
[Service]
|
||||||
|
Type=oneshot
|
||||||
|
ExecStart=/usr/bin/python3 /home/super/Projects/NetVM/bin/box-work.py loopback
|
||||||
|
WorkingDirectory=/home/super/Projects/NetVM
|
||||||
|
StandardOutput=journal
|
||||||
|
StandardError=journal
|
||||||
@@ -0,0 +1,10 @@
|
|||||||
|
[Unit]
|
||||||
|
Description=Run Box Work Cognitive True Loopback Sweeper every 5 minutes
|
||||||
|
|
||||||
|
[Timer]
|
||||||
|
OnBootSec=1min
|
||||||
|
OnUnitActiveSec=5min
|
||||||
|
Persistent=true
|
||||||
|
|
||||||
|
[Install]
|
||||||
|
WantedBy=timers.target
|
||||||
@@ -0,0 +1,121 @@
|
|||||||
|
"""test_box_readback_loopback.py — Unit tests for Readback Gate and True Loopback Engine."""
|
||||||
|
|
||||||
|
import unittest
|
||||||
|
from unittest.mock import patch, MagicMock
|
||||||
|
from datetime import datetime, timezone, timedelta
|
||||||
|
import os
|
||||||
|
import sys
|
||||||
|
|
||||||
|
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", "bin"))
|
||||||
|
import box_readback_loopback as brl
|
||||||
|
|
||||||
|
|
||||||
|
class TestReadbackValidation(unittest.TestCase):
|
||||||
|
|
||||||
|
def test_strict_tag_match(self):
|
||||||
|
text = "[READBACK] Ticket #219 | Branch: dev/pip/219-work | Plan: run pytest"
|
||||||
|
valid, excerpt = brl.validate_readback(text, 219, "pip")
|
||||||
|
self.assertTrue(valid)
|
||||||
|
self.assertIn("Ticket #219", excerpt)
|
||||||
|
|
||||||
|
def test_semantic_fallback_match(self):
|
||||||
|
text = "Understood. I am working on ticket 219 on dev/pip/219-work."
|
||||||
|
valid, excerpt = brl.validate_readback(text, 219, "pip")
|
||||||
|
self.assertTrue(valid)
|
||||||
|
self.assertIn("219", excerpt)
|
||||||
|
|
||||||
|
def test_irrelevant_message_rejected(self):
|
||||||
|
text = "Heartbeat check-in completed, all systems green."
|
||||||
|
valid, excerpt = brl.validate_readback(text, 219, "pip")
|
||||||
|
self.assertFalse(valid)
|
||||||
|
self.assertEqual(excerpt, "")
|
||||||
|
|
||||||
|
def test_wrong_ticket_rejected(self):
|
||||||
|
text = "[READBACK] Ticket #218 | Branch: dev/pip/218-work"
|
||||||
|
valid, excerpt = brl.validate_readback(text, 219, "pip")
|
||||||
|
self.assertFalse(valid)
|
||||||
|
|
||||||
|
|
||||||
|
class TestLoopbackEscalation(unittest.TestCase):
|
||||||
|
|
||||||
|
@patch("box_readback_loopback.acp.get_passive_cognitive_state")
|
||||||
|
@patch("box_readback_loopback.gitea_api")
|
||||||
|
def test_loopback_silence_when_active(self, mock_gitea, mock_cog):
|
||||||
|
# 5 minutes inactive -> silence
|
||||||
|
now = datetime.now(timezone.utc)
|
||||||
|
recent = (now - timedelta(minutes=5)).isoformat()
|
||||||
|
mock_gitea.return_value = [
|
||||||
|
{"number": 219, "title": "Test", "created_at": recent, "assignee": {"username": "dev"}}
|
||||||
|
]
|
||||||
|
mock_cog.return_value = {"status": "IDLE", "cognitive_lock": False}
|
||||||
|
|
||||||
|
actions = brl.run_loopback_sweep(dry_run=True, verbose=False)
|
||||||
|
self.assertEqual(len(actions), 0)
|
||||||
|
|
||||||
|
@patch("box_readback_loopback.acp.get_passive_cognitive_state")
|
||||||
|
@patch("box_readback_loopback.gitea_api")
|
||||||
|
def test_loopback_tier1_after_15m_when_idle(self, mock_gitea, mock_cog):
|
||||||
|
# 20 minutes inactive -> Tier 1
|
||||||
|
now = datetime.now(timezone.utc)
|
||||||
|
old = (now - timedelta(minutes=20)).isoformat()
|
||||||
|
mock_gitea.side_effect = lambda ep, **kwargs: (
|
||||||
|
[{"number": 999, "title": "Stalled", "created_at": old, "assignee": {"username": "dev"}}]
|
||||||
|
if ep == "/repos/super/box/issues?state=open"
|
||||||
|
else []
|
||||||
|
)
|
||||||
|
mock_cog.return_value = {"status": "IDLE", "cognitive_lock": False}
|
||||||
|
|
||||||
|
actions = brl.run_loopback_sweep(dry_run=True, verbose=False)
|
||||||
|
self.assertTrue(any("Tier 1" in a for a in actions))
|
||||||
|
|
||||||
|
@patch("box_readback_loopback.acp.get_passive_cognitive_state")
|
||||||
|
@patch("box_readback_loopback.gitea_api")
|
||||||
|
def test_loopback_defers_when_thinking(self, mock_gitea, mock_cog):
|
||||||
|
# 25 minutes inactive, but agent is THINKING -> defer
|
||||||
|
now = datetime.now(timezone.utc)
|
||||||
|
old = (now - timedelta(minutes=25)).isoformat()
|
||||||
|
mock_gitea.side_effect = lambda ep, **kwargs: (
|
||||||
|
[{"number": 999, "title": "Stalled", "created_at": old, "assignee": {"username": "dev"}}]
|
||||||
|
if ep == "/repos/super/box/issues?state=open"
|
||||||
|
else []
|
||||||
|
)
|
||||||
|
mock_cog.return_value = {"status": "THINKING", "cognitive_lock": True}
|
||||||
|
|
||||||
|
actions = brl.run_loopback_sweep(dry_run=True, verbose=False)
|
||||||
|
self.assertEqual(len(actions), 0)
|
||||||
|
|
||||||
|
@patch("box_readback_loopback.acp.get_passive_cognitive_state")
|
||||||
|
@patch("box_readback_loopback.gitea_api")
|
||||||
|
def test_loopback_tier2_after_45m(self, mock_gitea, mock_cog):
|
||||||
|
# 50 minutes inactive -> Tier 2
|
||||||
|
now = datetime.now(timezone.utc)
|
||||||
|
old = (now - timedelta(minutes=50)).isoformat()
|
||||||
|
mock_gitea.side_effect = lambda ep, **kwargs: (
|
||||||
|
[{"number": 999, "title": "Stalled", "created_at": old, "assignee": {"username": "dev"}}]
|
||||||
|
if ep == "/repos/super/box/issues?state=open"
|
||||||
|
else []
|
||||||
|
)
|
||||||
|
mock_cog.return_value = {"status": "IDLE", "cognitive_lock": False}
|
||||||
|
|
||||||
|
actions = brl.run_loopback_sweep(dry_run=True, verbose=False)
|
||||||
|
self.assertTrue(any("Tier 2" in a for a in actions))
|
||||||
|
|
||||||
|
@patch("box_readback_loopback.acp.get_passive_cognitive_state")
|
||||||
|
@patch("box_readback_loopback.gitea_api")
|
||||||
|
def test_loopback_tier3_after_90m(self, mock_gitea, mock_cog):
|
||||||
|
# 95 minutes inactive -> Tier 3
|
||||||
|
now = datetime.now(timezone.utc)
|
||||||
|
old = (now - timedelta(minutes=95)).isoformat()
|
||||||
|
mock_gitea.side_effect = lambda ep, **kwargs: (
|
||||||
|
[{"number": 999, "title": "Stalled", "created_at": old, "assignee": {"username": "dev"}}]
|
||||||
|
if ep == "/repos/super/box/issues?state=open"
|
||||||
|
else []
|
||||||
|
)
|
||||||
|
mock_cog.return_value = {"status": "IDLE", "cognitive_lock": False}
|
||||||
|
|
||||||
|
actions = brl.run_loopback_sweep(dry_run=True, verbose=False)
|
||||||
|
self.assertTrue(any("Tier 3" in a for a in actions))
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
unittest.main()
|
||||||
Reference in New Issue
Block a user