1612 lines
63 KiB
Python
Executable File
1612 lines
63 KiB
Python
Executable File
#!/usr/bin/env python3
|
|
"""box-fleet-tui.py — Fleet control TUI for NetVM (stdlib curses).
|
|
|
|
Surfaces 6 views:
|
|
[1] Timers — user timer states (NEXT/LAST per fleet timer).
|
|
[2] Harvest — harvest watermarks per agent with freshness highlight.
|
|
[3] Followups — follow-up queue counts by status.
|
|
[4] Approvals — pending approvals across the fleet.
|
|
[5] Activity — per-agent last activity.
|
|
[6] Runtimes — agent sockets/sessions vs the fleet manifest, task
|
|
queue depths, and the reconcile dry-run plan.
|
|
ACTIONABLE: R runs the reconciler now, o/Enter
|
|
attaches to the selected session. Tabs 1-5 stay
|
|
read-only.
|
|
|
|
All data-gathering lives in pure, testable functions taking injected
|
|
runners/readers (see gather_*). The curses UI is a thin renderer over
|
|
those functions. Missing files/commands yield "n/a" rows — never a crash.
|
|
|
|
Usage:
|
|
box fleet-tui (wired in bin/super-cli.py)
|
|
box tui fleet
|
|
python3 bin/box-fleet-tui.py [--once [--json]] # non-interactive dump
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import curses
|
|
import json
|
|
import os
|
|
import subprocess
|
|
import sys
|
|
import time
|
|
from datetime import datetime, timezone
|
|
from pathlib import Path
|
|
from typing import Any, Callable, Dict, Iterable, List, Optional, Tuple
|
|
import urllib.request
|
|
|
|
try:
|
|
import agent_cognitive_probe as acp
|
|
except ImportError:
|
|
acp = None
|
|
|
|
REPO_ROOT = Path(__file__).resolve().parent.parent
|
|
BIN_DIR = REPO_ROOT / "bin"
|
|
if str(BIN_DIR) not in sys.path:
|
|
sys.path.insert(0, str(BIN_DIR))
|
|
SYSTEMD_DIR = REPO_ROOT / "systemd"
|
|
WATERMARKS_FILE = REPO_ROOT / "siphon-watermarks.json"
|
|
FOLLOWUPS_FILE = REPO_ROOT / "followups.json"
|
|
JOB_SIDECHATS_FILE = REPO_ROOT / "job-sidechats.json"
|
|
DM_LOG = REPO_ROOT / "dm-log.jsonl"
|
|
JOB_LOG = REPO_ROOT / "job-log.jsonl"
|
|
CHAT_HISTORY_LOG = REPO_ROOT / "logs" / "chat-history.jsonl"
|
|
FLEET_MANIFEST = REPO_ROOT / "fleet" / "agents.json"
|
|
FLEET_TASKS_DIR = REPO_ROOT / "fleet" / "tasks"
|
|
FLEET_DONE_DIR = FLEET_TASKS_DIR / "done"
|
|
|
|
FLEET_AGENTS = ["muse", "pip", "646", "opm", "def", "dev"]
|
|
|
|
NA = "n/a"
|
|
|
|
# Freshness thresholds (seconds) for harvest watermarks.
|
|
FRESH_S = 15 * 60
|
|
AGING_S = 2 * 60 * 60
|
|
|
|
RunFn = Callable[..., Tuple[int, str]]
|
|
ReadJsonFn = Callable[[Path], Optional[Any]]
|
|
IterLinesFn = Callable[[Path], Iterable[str]]
|
|
|
|
EPOCH = datetime(1970, 1, 1, tzinfo=timezone.utc)
|
|
|
|
|
|
def _safe_str(val: Any, default: str = "") -> str:
|
|
"""Safe string converter that turns None into default."""
|
|
return default if val is None else str(val)
|
|
|
|
|
|
|
|
# =====================================================================
|
|
# Default IO primitives (injectable seams for tests)
|
|
# =====================================================================
|
|
|
|
def _run(cmd: List[str], timeout: int = 15) -> Tuple[int, str]:
|
|
"""Run cmd, capture output. Returns (returncode, combined_output)."""
|
|
try:
|
|
r = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout)
|
|
return r.returncode, ((r.stdout or "") + (r.stderr or "")).strip()
|
|
except subprocess.TimeoutExpired:
|
|
return 124, "timed out after %ds: %s" % (timeout, " ".join(cmd))
|
|
except OSError as e:
|
|
return 127, str(e)
|
|
|
|
|
|
def _read_json(path: Path) -> Optional[Any]:
|
|
"""Read + parse a JSON file. Missing/corrupt -> None (never raises)."""
|
|
try:
|
|
with open(path, "r") as f:
|
|
return json.load(f)
|
|
except Exception:
|
|
return None
|
|
|
|
|
|
def _iter_lines(path: Path, max_bytes: int = 65536) -> Iterable[str]:
|
|
"""Yield trailing lines of a (possibly large) text file, oldest first.
|
|
|
|
Reads only the last max_bytes. Missing/unreadable -> yields nothing.
|
|
"""
|
|
try:
|
|
with open(path, "rb") as f:
|
|
f.seek(0, 2)
|
|
size = f.tell()
|
|
f.seek(max(0, size - max_bytes))
|
|
chunk = f.read().decode("utf-8", errors="replace")
|
|
except Exception:
|
|
return
|
|
lines = chunk.splitlines()
|
|
# Drop a probable partial first line when we sliced mid-file.
|
|
if size > max_bytes and lines:
|
|
lines = lines[1:]
|
|
for line in lines:
|
|
yield line
|
|
|
|
|
|
# =====================================================================
|
|
# Small pure helpers
|
|
# =====================================================================
|
|
|
|
def parse_ts(value: Any) -> Optional[datetime]:
|
|
"""Parse an ISO-8601 timestamp (or epoch number) into aware datetime."""
|
|
if value is None:
|
|
return None
|
|
if isinstance(value, (int, float)):
|
|
try:
|
|
return datetime.fromtimestamp(float(value), tz=timezone.utc)
|
|
except Exception:
|
|
return None
|
|
if not isinstance(value, str) or not value.strip():
|
|
return None
|
|
try:
|
|
dt = datetime.fromisoformat(value.strip().replace("Z", "+00:00"))
|
|
if dt.tzinfo is None:
|
|
dt = dt.replace(tzinfo=timezone.utc)
|
|
return dt
|
|
except Exception:
|
|
return None
|
|
|
|
|
|
def relative_age(value: Any, now: Optional[datetime] = None) -> str:
|
|
"""Short relative age ("5m ago") for a timestamp; "n/a" when unknown."""
|
|
dt = parse_ts(value)
|
|
if dt is None:
|
|
return NA
|
|
now = now or datetime.now(timezone.utc)
|
|
diff = int((now - dt).total_seconds())
|
|
if diff < 0:
|
|
return "just now"
|
|
if diff < 60:
|
|
return "%ds ago" % diff
|
|
if diff < 3600:
|
|
return "%dm ago" % (diff // 60)
|
|
if diff < 86400:
|
|
return "%dh ago" % (diff // 3600)
|
|
return "%dd ago" % (diff // 86400)
|
|
|
|
|
|
def freshness_bucket(value: Any, now: Optional[datetime] = None) -> str:
|
|
"""Freshness of a timestamp: fresh / aging / stale / unknown."""
|
|
dt = parse_ts(value)
|
|
if dt is None:
|
|
return "unknown"
|
|
now = now or datetime.now(timezone.utc)
|
|
age = (now - dt).total_seconds()
|
|
if age < 0:
|
|
return "fresh"
|
|
if age < FRESH_S:
|
|
return "fresh"
|
|
if age < AGING_S:
|
|
return "aging"
|
|
return "stale"
|
|
|
|
|
|
def short_ts(value: Any) -> str:
|
|
"""Compact "MM-DD HH:MM" display for a timestamp; "n/a" when unknown."""
|
|
dt = parse_ts(value)
|
|
if dt is None:
|
|
return NA
|
|
try:
|
|
local = dt.astimezone()
|
|
except Exception:
|
|
local = dt
|
|
return local.strftime("%m-%d %H:%M")
|
|
|
|
|
|
# =====================================================================
|
|
# Surface 1: user timer states (NEXT/LAST per fleet timer)
|
|
# =====================================================================
|
|
|
|
def shipped_timers(systemd_dir: Optional[Path] = None,
|
|
listdir: Optional[Callable[[Path], List[str]]] = None) -> List[str]:
|
|
"""Sorted *.timer unit names shipped in systemd/. Missing dir -> []."""
|
|
d = Path(systemd_dir) if systemd_dir else SYSTEMD_DIR
|
|
try:
|
|
if listdir is not None:
|
|
names = listdir(d)
|
|
else:
|
|
names = [p.name for p in d.iterdir() if p.is_file()]
|
|
except Exception:
|
|
return []
|
|
return sorted(n for n in names if n.endswith(".timer"))
|
|
|
|
|
|
def _parse_list_timers_json(out: str) -> Dict[str, Dict[str, Any]]:
|
|
"""Parse `systemctl list-timers --output=json` (µs-since-epoch ints)."""
|
|
data = json.loads(out)
|
|
states: Dict[str, Dict[str, Any]] = {}
|
|
for entry in data:
|
|
if not isinstance(entry, dict):
|
|
continue
|
|
unit = entry.get("unit")
|
|
if not unit:
|
|
continue
|
|
nxt = entry.get("next") or 0
|
|
last = entry.get("last") or 0
|
|
try:
|
|
next_ts = float(nxt) / 1e6 if float(nxt) > 0 else None
|
|
except Exception:
|
|
next_ts = None
|
|
try:
|
|
last_ts = float(last) / 1e6 if float(last) > 0 else None
|
|
except Exception:
|
|
last_ts = None
|
|
states[unit] = {
|
|
"next_ts": next_ts,
|
|
"last_ts": last_ts,
|
|
"activates": entry.get("activates") or NA,
|
|
}
|
|
return states
|
|
|
|
|
|
def _parse_list_timers_text(out: str) -> Dict[str, Dict[str, Any]]:
|
|
"""Fallback parser for the human-readable list-timers table.
|
|
|
|
NEXT/LAST columns carry weekday + date + time + zone, e.g.
|
|
"Wed 2026-10-07 05:37:00 UTC"; "-" marks an empty cell.
|
|
"""
|
|
states: Dict[str, Dict[str, Any]] = {}
|
|
for line in out.splitlines():
|
|
line = line.rstrip()
|
|
if not line or line.startswith("NEXT"):
|
|
continue
|
|
parts = line.split()
|
|
if len(parts) < 3:
|
|
continue
|
|
unit = parts[-2]
|
|
activates = parts[-1]
|
|
if not unit.endswith(".timer"):
|
|
continue
|
|
# Walk cells: NEXT is 4 tokens ("Wed DATE TIME TZ") or "-",
|
|
# then LEFT (1-3 tokens ending in "ago" or a duration), then LAST.
|
|
rest = parts[:-2]
|
|
try:
|
|
if rest and rest[0] == "-":
|
|
next_ts = None
|
|
rest = rest[1:]
|
|
else:
|
|
next_ts = _table_cell_ts(rest[:4])
|
|
rest = rest[4:]
|
|
# Skip LEFT cell: consume until LAST cell starts (weekday or "-").
|
|
weekdays = {"Mon", "Tue", "Wed", "Thu", "Fri", "Sat", "Sun"}
|
|
while rest and rest[0] not in weekdays and rest[0] != "-":
|
|
rest = rest[1:]
|
|
if rest and rest[0] == "-":
|
|
last_ts = None
|
|
else:
|
|
last_ts = _table_cell_ts(rest[:4])
|
|
except Exception:
|
|
next_ts, last_ts = None, None
|
|
states[unit] = {"next_ts": next_ts, "last_ts": last_ts,
|
|
"activates": activates}
|
|
return states
|
|
|
|
|
|
def _table_cell_ts(tokens: List[str]) -> Optional[float]:
|
|
if len(tokens) < 4:
|
|
return None
|
|
try:
|
|
dt = datetime.strptime(" ".join(tokens[1:3]), "%Y-%m-%d %H:%M:%S")
|
|
dt = dt.replace(tzinfo=timezone.utc)
|
|
return (dt - EPOCH).total_seconds()
|
|
except Exception:
|
|
return None
|
|
|
|
|
|
def gather_timers(run: Optional[RunFn] = None,
|
|
timer_names: Optional[List[str]] = None,
|
|
now: Optional[datetime] = None) -> Dict[str, Any]:
|
|
"""NEXT/LAST state per fleet timer.runner failure -> "n/a" rows.
|
|
|
|
Returns {"rows": [{timer, next, next_rel, last, last_rel, activates,
|
|
present}], "note": str}. `timer_names=None` reports every non-transient
|
|
timer the manager lists (run-*.timer excluded).
|
|
"""
|
|
run = run or _run
|
|
now = now or datetime.now(timezone.utc)
|
|
states: Dict[str, Dict[str, Any]] = {}
|
|
note = ""
|
|
try:
|
|
rc, out = run(["systemctl", "--user", "list-timers", "--no-pager",
|
|
"--all", "--output=json"], timeout=15)
|
|
except Exception as e:
|
|
rc, out = 127, str(e)
|
|
if rc == 0 and out.strip():
|
|
try:
|
|
states = _parse_list_timers_json(out)
|
|
except Exception:
|
|
note = "json parse failed; tried text fallback. "
|
|
try:
|
|
rc2, out2 = run(["systemctl", "--user", "list-timers",
|
|
"--no-pager", "--all"], timeout=15)
|
|
if rc2 == 0:
|
|
states = _parse_list_timers_text(out2)
|
|
except Exception:
|
|
states = {}
|
|
else:
|
|
note = "systemctl unavailable (%s). " % (out[:60] if out else "rc=%s" % rc)
|
|
|
|
if timer_names is None:
|
|
names = sorted(u for u in states if not u.startswith("run-"))
|
|
else:
|
|
names = list(timer_names)
|
|
rows = []
|
|
for name in names:
|
|
st = states.get(name)
|
|
if st is None:
|
|
rows.append({"timer": name, "next": NA, "next_rel": NA,
|
|
"last": NA, "last_rel": NA, "activates": NA,
|
|
"present": False})
|
|
continue
|
|
nts, lts = st.get("next_ts"), st.get("last_ts")
|
|
rows.append({
|
|
"timer": name,
|
|
"next": short_ts(nts) if nts else NA,
|
|
"next_rel": _future_rel(nts, now),
|
|
"last": short_ts(lts) if lts else NA,
|
|
"last_rel": relative_age(lts, now) if lts else NA,
|
|
"activates": st.get("activates") or NA,
|
|
"present": True,
|
|
})
|
|
return {"rows": rows, "note": note.strip()}
|
|
|
|
|
|
def _future_rel(ts: Optional[float], now: datetime) -> str:
|
|
if not ts:
|
|
return NA
|
|
diff = int(ts - now.timestamp())
|
|
if diff < 0:
|
|
return "due"
|
|
if diff < 60:
|
|
return "in %ds" % diff
|
|
if diff < 3600:
|
|
return "in %dm" % (diff // 60)
|
|
if diff < 86400:
|
|
return "in %dh" % (diff // 3600)
|
|
return "in %dd" % (diff // 86400)
|
|
|
|
|
|
# =====================================================================
|
|
# Surface 2: harvest watermarks per agent (freshness highlighted)
|
|
# =====================================================================
|
|
|
|
def _latest_harvest_ts(iter_lines: IterLinesFn) -> Dict[str, str]:
|
|
"""Latest chat-history ts per "agent:thread_id" key."""
|
|
latest: Dict[str, str] = {}
|
|
try:
|
|
lines = iter_lines(CHAT_HISTORY_LOG)
|
|
except Exception:
|
|
return latest
|
|
try:
|
|
for line in lines:
|
|
if not line or not line.strip():
|
|
continue
|
|
try:
|
|
rec = json.loads(line)
|
|
except Exception:
|
|
continue
|
|
if not isinstance(rec, dict):
|
|
continue
|
|
agent, tid = rec.get("agent"), rec.get("thread_id")
|
|
ts = rec.get("ts")
|
|
if agent and tid and ts:
|
|
latest["%s:%s" % (agent, tid)] = ts
|
|
except Exception:
|
|
pass
|
|
return latest
|
|
|
|
|
|
def gather_harvest(read_json: Optional[ReadJsonFn] = None,
|
|
iter_lines: Optional[IterLinesFn] = None,
|
|
agents: Optional[List[str]] = None,
|
|
now: Optional[datetime] = None) -> Dict[str, Any]:
|
|
"""Watermark + freshness rows per agent main chat and sidechat threads.
|
|
|
|
Returns {"rows": [{agent, thread_name, thread_id, watermark, last,
|
|
last_rel, freshness}], "note": str}. Missing files -> "n/a" cells.
|
|
"""
|
|
read_json = read_json or _read_json
|
|
iter_lines = iter_lines or _iter_lines
|
|
agents = agents or FLEET_AGENTS
|
|
now = now or datetime.now(timezone.utc)
|
|
notes: List[str] = []
|
|
|
|
watermarks: Dict[str, Any] = {}
|
|
try:
|
|
loaded = read_json(WATERMARKS_FILE)
|
|
if isinstance(loaded, dict):
|
|
watermarks = loaded
|
|
else:
|
|
notes.append("watermarks file missing/empty.")
|
|
except Exception:
|
|
notes.append("watermarks unreadable.")
|
|
|
|
sidechats: Dict[str, Any] = {}
|
|
try:
|
|
loaded = read_json(JOB_SIDECHATS_FILE)
|
|
if isinstance(loaded, dict):
|
|
sidechats = loaded
|
|
except Exception:
|
|
pass
|
|
|
|
latest = _latest_harvest_ts(iter_lines)
|
|
|
|
# Sidechat threads grouped by owning agent.
|
|
owned: Dict[str, List[Tuple[str, str]]] = {a: [] for a in agents}
|
|
for alias, val in sidechats.items():
|
|
if alias.startswith("_") or not isinstance(val, dict):
|
|
continue
|
|
tid = val.get("thread_uuid") or val.get("uuid")
|
|
owner = val.get("agent", "opm")
|
|
if tid and owner in owned:
|
|
owned[owner].append((alias, tid))
|
|
|
|
rows = []
|
|
for agent in agents:
|
|
targets: List[Tuple[str, str]] = [("Main Chat", "main")]
|
|
targets.extend(sorted(owned.get(agent, [])))
|
|
for thread_name, tid in targets:
|
|
key = "%s:%s" % (agent, tid)
|
|
wm = watermarks.get(key)
|
|
wm_s = str(wm)[:24] if wm else NA
|
|
last_ts = latest.get(key)
|
|
rows.append({
|
|
"agent": agent,
|
|
"thread_name": thread_name,
|
|
"thread_id": tid,
|
|
"watermark": wm_s,
|
|
"last": short_ts(last_ts) if last_ts else NA,
|
|
"last_rel": relative_age(last_ts, now) if last_ts else NA,
|
|
"freshness": freshness_bucket(last_ts, now) if last_ts else "unknown",
|
|
})
|
|
return {"rows": rows, "note": " ".join(notes)}
|
|
|
|
|
|
# =====================================================================
|
|
# Surface 3: follow-up queue counts by status
|
|
# =====================================================================
|
|
|
|
FOLLOWUP_STATUSES = ("pending", "escalated", "resolved")
|
|
|
|
|
|
def gather_followups(read_json: Optional[ReadJsonFn] = None,
|
|
now: Optional[datetime] = None) -> Dict[str, Any]:
|
|
"""Queue counts by status + oldest pending item. Corrupt -> zeros.
|
|
|
|
Returns {"counts": {status: n}, "total": n, "overdue": n,
|
|
"oldest_pending": {...}|None, "by_recipient": {agent: n}, "note": str}.
|
|
"""
|
|
read_json = read_json or _read_json
|
|
now = now or datetime.now(timezone.utc)
|
|
try:
|
|
loaded = read_json(FOLLOWUPS_FILE)
|
|
followups = loaded if isinstance(loaded, dict) else {}
|
|
note = "" if loaded is not None else "followups file missing/empty."
|
|
except Exception:
|
|
followups, note = {}, "followups unreadable."
|
|
|
|
counts: Dict[str, int] = {s: 0 for s in FOLLOWUP_STATUSES}
|
|
other = 0
|
|
overdue = 0
|
|
by_recipient: Dict[str, int] = {}
|
|
oldest: Optional[Dict[str, Any]] = None
|
|
oldest_dt: Optional[datetime] = None
|
|
for dm_id, rec in followups.items():
|
|
if not isinstance(rec, dict):
|
|
other += 1
|
|
continue
|
|
st = str(rec.get("status", "pending"))
|
|
if st in counts:
|
|
counts[st] += 1
|
|
else:
|
|
other += 1
|
|
if st == "pending":
|
|
recip = str(rec.get("recipient", "?"))
|
|
by_recipient[recip] = by_recipient.get(recip, 0) + 1
|
|
dl = parse_ts(rec.get("deadline"))
|
|
if dl is not None and dl < now:
|
|
overdue += 1
|
|
sent = parse_ts(rec.get("sent_at"))
|
|
if sent is not None and (oldest_dt is None or sent < oldest_dt):
|
|
oldest_dt = sent
|
|
oldest = {
|
|
"dm_id": str(dm_id)[:8],
|
|
"route": "%s -> %s" % (rec.get("sender", "?"),
|
|
rec.get("recipient", "?")),
|
|
"target": str(rec.get("target", "main"))[:24],
|
|
"sent": short_ts(rec.get("sent_at")),
|
|
"sent_rel": relative_age(rec.get("sent_at"), now),
|
|
}
|
|
if other:
|
|
counts["other"] = other
|
|
total = sum(counts.values())
|
|
return {"counts": counts, "total": total, "overdue": overdue,
|
|
"oldest_pending": oldest, "by_recipient": by_recipient,
|
|
"note": note}
|
|
|
|
|
|
# =====================================================================
|
|
# Surface 4: pending approvals
|
|
# =====================================================================
|
|
|
|
def _default_approval_check(nodes: List[str]) -> List[Dict[str, Any]]:
|
|
"""Lazy approvals.check_fleet_approvals with per-node isolation."""
|
|
try:
|
|
sys.path.insert(0, str(BIN_DIR))
|
|
import approvals
|
|
except Exception as e:
|
|
raise RuntimeError("approvals module unavailable: %s" % e)
|
|
results = []
|
|
for node in nodes:
|
|
try:
|
|
results.append(approvals.inspect_node_approvals(node))
|
|
except Exception as e:
|
|
results.append({"node": node, "status": "ERROR",
|
|
"has_pending": False, "error": str(e)})
|
|
return results
|
|
|
|
|
|
def gather_approvals(check: Optional[Callable[[List[str]], List[Dict[str, Any]]]] = None,
|
|
nodes: Optional[List[str]] = None) -> Dict[str, Any]:
|
|
"""Pending approvals across the fleet. Probe failure -> note, no crash.
|
|
|
|
Returns {"pending": [{node, status, title, target, trusted}],
|
|
"total_pending": n, "unreachable": [nodes], "checked": n, "note": str}.
|
|
"""
|
|
nodes = nodes or FLEET_AGENTS
|
|
check = check or _default_approval_check
|
|
try:
|
|
results = check(list(nodes))
|
|
except Exception as e:
|
|
return {"pending": [], "total_pending": 0, "unreachable": [],
|
|
"checked": 0, "note": "approval probe failed: %s" % e}
|
|
if not isinstance(results, list):
|
|
return {"pending": [], "total_pending": 0, "unreachable": [],
|
|
"checked": 0, "note": "approval probe returned no data."}
|
|
pending, unreachable = [], []
|
|
for item in results:
|
|
if not isinstance(item, dict):
|
|
continue
|
|
node = str(item.get("node", "?"))
|
|
status = str(item.get("status", "UNKNOWN"))
|
|
if status == "UNREACHABLE":
|
|
unreachable.append(node)
|
|
if item.get("has_pending"):
|
|
title = str(item.get("title", ""))[:60] or status
|
|
pending.append({
|
|
"node": node,
|
|
"status": status,
|
|
"title": title,
|
|
"target": str(item.get("target", "-"))[:32],
|
|
"trusted": bool(item.get("is_trusted", False)),
|
|
})
|
|
return {"pending": pending, "total_pending": len(pending),
|
|
"unreachable": unreachable, "checked": len(results), "note": ""}
|
|
|
|
|
|
# =====================================================================
|
|
# Surface 5: per-agent last activity
|
|
# =====================================================================
|
|
|
|
def _scan_activity(iter_lines: IterLinesFn) -> Dict[str, Dict[str, str]]:
|
|
"""Last-seen {agent: {ts, source, detail}} across dm/job/chat logs."""
|
|
last: Dict[str, Dict[str, str]] = {}
|
|
sources: List[Tuple[Path, str]] = [
|
|
(DM_LOG, "dm"), (JOB_LOG, "job"), (CHAT_HISTORY_LOG, "chat"),
|
|
]
|
|
for path, source in sources:
|
|
try:
|
|
lines = iter_lines(path)
|
|
except Exception:
|
|
continue
|
|
try:
|
|
for line in lines:
|
|
if not line or not line.strip():
|
|
continue
|
|
try:
|
|
rec = json.loads(line)
|
|
except Exception:
|
|
continue
|
|
if not isinstance(rec, dict):
|
|
continue
|
|
agent = rec.get("agent")
|
|
ts = rec.get("ts")
|
|
if not agent or not ts or parse_ts(ts) is None:
|
|
continue
|
|
prev = last.get(str(agent))
|
|
if prev is not None and str(prev["ts"]) >= str(ts):
|
|
continue
|
|
if source == "dm":
|
|
detail = "%s -> %s (%s)" % (
|
|
rec.get("agent", "?"), rec.get("to", "?"),
|
|
rec.get("target", rec.get("type", "?")))
|
|
elif source == "job":
|
|
detail = str(rec.get("result_snippet", rec.get("type", "?")))[:48]
|
|
else:
|
|
detail = str(rec.get("thread_name", rec.get("thread_id", "?")))[:40]
|
|
last[str(agent)] = {"ts": str(ts), "source": source,
|
|
"detail": detail}
|
|
except Exception:
|
|
continue
|
|
return last
|
|
|
|
|
|
def gather_activity(iter_lines: Optional[IterLinesFn] = None,
|
|
agents: Optional[List[str]] = None,
|
|
now: Optional[datetime] = None) -> Dict[str, Any]:
|
|
"""Per-agent last activity rows. Empty logs -> "n/a" rows, never crash.
|
|
|
|
Returns {"rows": [{agent, last, last_rel, source, detail, freshness}],
|
|
"note": str}.
|
|
"""
|
|
iter_lines = iter_lines or _iter_lines
|
|
agents = agents or FLEET_AGENTS
|
|
now = now or datetime.now(timezone.utc)
|
|
try:
|
|
last = _scan_activity(iter_lines)
|
|
except Exception as e:
|
|
last = {}
|
|
note = "activity scan failed: %s" % e
|
|
else:
|
|
note = "" if last else "no activity records found."
|
|
rows = []
|
|
for agent in agents:
|
|
rec = last.get(agent)
|
|
if rec is None:
|
|
rows.append({"agent": agent, "last": NA, "last_rel": NA,
|
|
"source": NA, "detail": NA, "freshness": "unknown"})
|
|
else:
|
|
rows.append({
|
|
"agent": agent,
|
|
"last": short_ts(rec["ts"]),
|
|
"last_rel": relative_age(rec["ts"], now),
|
|
"source": rec["source"],
|
|
"detail": rec["detail"],
|
|
"freshness": freshness_bucket(rec["ts"], now),
|
|
})
|
|
return {"rows": rows, "note": note}
|
|
|
|
|
|
def gather_runtimes(rows_fn: Optional[Callable[[], List[Dict[str, Any]]]] = None,
|
|
manifest_fn: Optional[Callable[[], Dict[str, Any]]] = None,
|
|
state_fn: Optional[Callable[[], Dict[str, Any]]] = None,
|
|
queue_fn: Optional[Callable[[], Dict[str, Any]]] = None,
|
|
done_count_fn: Optional[Callable[[], int]] = None,
|
|
brief_read: Optional[Callable[[str], Optional[str]]] = None,
|
|
plan_fn: Optional[Callable[[], Dict[str, Any]]] = None,
|
|
fleet_sockets: Optional[List[str]] = None,
|
|
now: Optional[datetime] = None) -> Dict[str, Any]:
|
|
"""Agent sockets/sessions vs manifest + queue + reconcile plan.
|
|
|
|
Every input is injectable; defaults hit the live fleet (tmux rows,
|
|
manifest, reconciler state, queue, dry-run plan). Never raises:
|
|
backend failures yield notes and n/a fields.
|
|
"""
|
|
try:
|
|
import muse_choice_watcher as mcw
|
|
import runtime_reconcile as rec
|
|
except Exception as e:
|
|
return {"note": "runtimes backend unavailable (%s)" % e, "agents": []}
|
|
|
|
if fleet_sockets is None:
|
|
try:
|
|
fleet_sockets = list(mcw.FLEET_SOCKETS)
|
|
except Exception:
|
|
fleet_sockets = []
|
|
manifest_path = str(FLEET_MANIFEST)
|
|
if manifest_fn is None:
|
|
manifest_fn = lambda: rec.load_manifest(manifest_path) # noqa: E731
|
|
if state_fn is None:
|
|
state_fn = lambda: rec.load_state( # noqa: E731
|
|
str(REPO_ROOT / "fleet" / "state.json"))
|
|
if queue_fn is None:
|
|
queue_fn = lambda: rec.queue_status(str(FLEET_TASKS_DIR)) # noqa: E731
|
|
if rows_fn is None:
|
|
rows_fn = mcw.all_runtime_rows # noqa: E731
|
|
if plan_fn is None:
|
|
plan_fn = lambda: rec.run_reconcile( # noqa: E731
|
|
manifest_path, dry_run=True)
|
|
if brief_read is None:
|
|
def brief_read(path: str) -> Optional[str]:
|
|
try:
|
|
with open(path) as f:
|
|
return f.read().strip()
|
|
except Exception:
|
|
return None
|
|
|
|
notes: List[str] = []
|
|
try:
|
|
manifest = manifest_fn()
|
|
except Exception as e:
|
|
manifest = {"agents": [], "errors": [str(e)]}
|
|
agents_declared = manifest.get("agents", []) or []
|
|
for err in manifest.get("errors", []) or []:
|
|
if "no manifest" not in str(err):
|
|
notes.append("manifest: %s" % err)
|
|
try:
|
|
rows = rows_fn() or []
|
|
except Exception as e:
|
|
rows = []
|
|
notes.append("runtime rows: %s" % e)
|
|
try:
|
|
state = state_fn() or {}
|
|
except Exception:
|
|
state = {}
|
|
records = state.get("agents") or {}
|
|
try:
|
|
queue = queue_fn() or {}
|
|
except Exception as e:
|
|
queue = {}
|
|
notes.append("queue: %s" % e)
|
|
if done_count_fn is None:
|
|
try:
|
|
done_names = [n for n in os.listdir(FLEET_DONE_DIR)
|
|
if not n.startswith(".")]
|
|
done_count = len(done_names)
|
|
except Exception:
|
|
done_count = 0
|
|
else:
|
|
try:
|
|
done_count = done_count_fn()
|
|
except Exception:
|
|
done_count = 0
|
|
|
|
live_by_session: Dict[str, List[Dict[str, Any]]] = {}
|
|
for row in rows:
|
|
live_by_session.setdefault(row.get("session", ""), []).append(row)
|
|
|
|
agents: List[Dict[str, Any]] = []
|
|
for entry in agents_declared:
|
|
session = entry.get("session", "?")
|
|
sock = entry.get("socket", "?")
|
|
panes = [r for r in live_by_session.get(session, [])
|
|
if r.get("socket") == sock]
|
|
pane = panes[0] if panes else None
|
|
key = rec.state_key(sock, session)
|
|
record = records.get(key)
|
|
briefed: Any = False
|
|
if record is not None:
|
|
text = brief_read(entry.get("brief", ""))
|
|
if text is None:
|
|
briefed = "unknown"
|
|
elif rec.brief_sha(text) == record.get("brief_sha"):
|
|
briefed = True
|
|
else:
|
|
briefed = "stale"
|
|
agents.append({
|
|
"socket": sock, "session": session,
|
|
"hat": entry.get("hat", "?"),
|
|
"enabled": bool(entry.get("enabled", True)),
|
|
"live": pane is not None,
|
|
"pane": pane.get("pane") if pane else NA,
|
|
"state": pane.get("state", NA) if pane else NA,
|
|
"mode": (pane.get("effective_mode")
|
|
or pane.get("permission_mode") or NA) if pane else NA,
|
|
"watcher": bool(pane.get("watcher_alive")) if pane else False,
|
|
"briefed": briefed,
|
|
})
|
|
|
|
manifest_sessions = {a["session"] for a in agents}
|
|
strays: List[Dict[str, Any]] = []
|
|
for row in rows:
|
|
if (row.get("socket") in (fleet_sockets or [])
|
|
and row.get("session") not in manifest_sessions):
|
|
strays.append({
|
|
"socket": row.get("socket", "?"),
|
|
"session": row.get("session", "?"),
|
|
"pane": row.get("pane", NA),
|
|
"cmd": (row.get("cmd") or "")[:24],
|
|
"state": row.get("state", NA),
|
|
})
|
|
|
|
try:
|
|
plan_report = plan_fn() or {}
|
|
except Exception as e:
|
|
plan_report = {}
|
|
notes.append("reconcile plan: %s" % e)
|
|
plan = {"launch": [], "brief": [], "nudge": [], "failed": []}
|
|
for item in plan_report.get("agents", []) or []:
|
|
action = item.get("action")
|
|
if action in plan:
|
|
plan[action].append(item.get("session", "?"))
|
|
for nudge in plan_report.get("nudges", []) or []:
|
|
plan["nudge"].append("%s:%s" % (nudge.get("session", "?"),
|
|
nudge.get("kind", "?")))
|
|
|
|
last_run: Any = NA
|
|
try:
|
|
ts = (state or {}).get("last_run")
|
|
if ts:
|
|
last_run = relative_age(float(ts), now)
|
|
except Exception:
|
|
last_run = NA
|
|
|
|
return {
|
|
"agents": agents,
|
|
"strays": strays,
|
|
"queue": {"pending": queue.get("pending", 0) if queue else 0,
|
|
"claimed": queue.get("claimed", {}) if queue else {},
|
|
"done": done_count},
|
|
"last_run": last_run,
|
|
"plan": plan,
|
|
"plan_errors": (plan_report.get("errors", []) or [])
|
|
+ (plan_report.get("claims", {}).get("errors", []) or []),
|
|
"note": "; ".join(notes),
|
|
}
|
|
|
|
|
|
def attach_argv(sock: str, session: str) -> List[str]:
|
|
"""tmux attach argv for a runtime open. Pure (tested)."""
|
|
return ["tmux", "-S", sock, "attach-session", "-t", session]
|
|
|
|
|
|
def do_reconcile_now(manifest_path: Optional[str] = None,
|
|
run_fn: Optional[Callable[..., Dict[str, Any]]] = None
|
|
) -> Tuple[bool, str]:
|
|
"""Run the runtime reconciler now, summarize the report.
|
|
|
|
Returns (ok, one-line summary). Never raises.
|
|
"""
|
|
try:
|
|
if run_fn is None:
|
|
import runtime_reconcile as rec
|
|
run_fn = rec.run_reconcile
|
|
report = run_fn(manifest_path or str(FLEET_MANIFEST))
|
|
except Exception as e:
|
|
return False, "reconcile failed: %s" % e
|
|
try:
|
|
actions = [a.get("action", "?")
|
|
for a in report.get("agents", []) or []]
|
|
launched = sum(1 for a in actions if a == "launched")
|
|
briefed = sum(1 for a in actions if a == "briefed")
|
|
failed = sum(1 for a in actions if a == "failed")
|
|
nudged = len(report.get("nudges", []) or [])
|
|
requeued = len((report.get("claims", {}) or {}).get("requeued", [])
|
|
or [])
|
|
ok_count = sum(1 for a in actions if a in ("ok", "adopted"))
|
|
msg = ("reconcile: %d ok, %d launched, %d briefed, "
|
|
"%d nudged, %d requeued" % (
|
|
ok_count, launched, briefed, nudged, requeued))
|
|
if failed:
|
|
msg += ", %d FAILED" % failed
|
|
return bool(report.get("ok", True)), msg
|
|
except Exception as e:
|
|
return False, "reconcile report unreadable: %s" % e
|
|
|
|
|
|
def gather_work(run: Optional[RunFn] = None,
|
|
read_json: Optional[ReadJsonFn] = None,
|
|
now: Optional[datetime] = None) -> Dict[str, Any]:
|
|
"""Gather real-time cognitive sensor states, Gitea build issues, and loopback state."""
|
|
_run_fn = run or _run
|
|
_read_fn = read_json or _read_json
|
|
|
|
cog_states = {}
|
|
if acp:
|
|
for agent in FLEET_AGENTS:
|
|
try:
|
|
cog_states[agent] = acp.get_passive_cognitive_state(agent)
|
|
except Exception:
|
|
cog_states[agent] = {"status": "UNKNOWN", "cognitive_lock": False, "badge": "⚠️ UNKNOWN"}
|
|
else:
|
|
for agent in FLEET_AGENTS:
|
|
cog_states[agent] = {"status": "NO_PROBE", "cognitive_lock": False, "badge": "n/a"}
|
|
|
|
# Fetch recent build tickets from Gitea
|
|
issues = []
|
|
try:
|
|
req = urllib.request.Request(
|
|
"https://tea.muse-dev.online/api/v1/repos/super/box/issues?state=all&limit=8",
|
|
headers={"User-Agent": "Box-Fleet-TUI/1.0"}
|
|
)
|
|
with urllib.request.urlopen(req, timeout=3.0) as r:
|
|
issues = json.loads(r.read().decode())
|
|
except Exception:
|
|
pass
|
|
|
|
# Check loopback timer state
|
|
rc, loop_txt = _run_fn(["systemctl", "--user", "is-active", "box-work-loopback.timer"], timeout=5)
|
|
loop_active = (rc == 0 and "active" in loop_txt.lower())
|
|
|
|
return {
|
|
"cognitive": cog_states,
|
|
"issues": issues if isinstance(issues, list) else [],
|
|
"loopback_active": loop_active,
|
|
}
|
|
|
|
|
|
def gather_all(run: Optional[RunFn] = None,
|
|
read_json: Optional[ReadJsonFn] = None,
|
|
iter_lines: Optional[IterLinesFn] = None,
|
|
approval_check: Optional[Callable] = None,
|
|
timer_names: Optional[List[str]] = None,
|
|
runtimes_kwargs: Optional[Dict[str, Any]] = None,
|
|
now: Optional[datetime] = None) -> Dict[str, Any]:
|
|
"""One-shot snapshot of all seven surfaces (TUI refresh + --once)."""
|
|
if timer_names is None:
|
|
timer_names = shipped_timers()
|
|
return {
|
|
"timers": gather_timers(run=run, timer_names=timer_names, now=now),
|
|
"harvest": gather_harvest(read_json=read_json, iter_lines=iter_lines,
|
|
now=now),
|
|
"followups": gather_followups(read_json=read_json, now=now),
|
|
"approvals": gather_approvals(check=approval_check),
|
|
"activity": gather_activity(iter_lines=iter_lines, now=now),
|
|
"runtimes": gather_runtimes(**(runtimes_kwargs or {}), now=now),
|
|
"work": gather_work(run=run, read_json=read_json, now=now),
|
|
}
|
|
|
|
|
|
# =====================================================================
|
|
# Curses UI (thin read-only renderer over gather_*)
|
|
# =====================================================================
|
|
|
|
AUTO_REFRESH_S = 60.0
|
|
|
|
|
|
class BoxFleetTUI:
|
|
"""Fleet control console. Tabs 1-5 read-only; 6:RUNTIMES acts."""
|
|
|
|
RUNTIMES_TAB = 5
|
|
WORK_TAB = 6
|
|
|
|
def __init__(self, stdscr: "curses.window"):
|
|
self.stdscr = stdscr
|
|
self.current_tab = 0
|
|
self.tabs = [
|
|
"1: TIMERS",
|
|
"2: HARVEST",
|
|
"3: FOLLOWUPS",
|
|
"4: APPROVALS",
|
|
"5: ACTIVITY",
|
|
"6: RUNTIMES",
|
|
"7: WORK",
|
|
]
|
|
self.rt_sel = 0
|
|
self.rt_items: List[Dict[str, str]] = []
|
|
self.work_sel = 0
|
|
try:
|
|
curses.curs_set(0)
|
|
except Exception:
|
|
pass
|
|
self.stdscr.nodelay(True)
|
|
self.stdscr.keypad(True)
|
|
if hasattr(curses, "set_escdelay"):
|
|
try:
|
|
curses.set_escdelay(25)
|
|
except Exception:
|
|
pass
|
|
self._init_colors()
|
|
|
|
self.timer_names = shipped_timers()
|
|
self.scroll = 0
|
|
self.show_help = False
|
|
self.status_msg = "Loading fleet snapshot..."
|
|
self.last_refresh = 0.0
|
|
self.snapshot: Dict[str, Any] = {}
|
|
self.refresh()
|
|
|
|
# -- setup ------------------------------------------------------
|
|
|
|
def _init_colors(self) -> None:
|
|
try:
|
|
curses.start_color()
|
|
curses.use_default_colors()
|
|
curses.init_pair(1, curses.COLOR_CYAN, -1)
|
|
curses.init_pair(2, curses.COLOR_YELLOW, -1)
|
|
curses.init_pair(3, curses.COLOR_GREEN, -1)
|
|
curses.init_pair(4, curses.COLOR_RED, -1)
|
|
curses.init_pair(5, curses.COLOR_MAGENTA, -1)
|
|
curses.init_pair(6, curses.COLOR_BLACK, curses.COLOR_CYAN)
|
|
curses.init_pair(7, curses.COLOR_BLACK, curses.COLOR_WHITE)
|
|
curses.init_pair(8, curses.COLOR_BLACK, curses.COLOR_YELLOW)
|
|
except Exception:
|
|
pass
|
|
|
|
def _attr(self, name: str) -> int:
|
|
try:
|
|
mapping = {
|
|
"normal": curses.color_pair(0),
|
|
"cyan": curses.color_pair(1) | curses.A_BOLD,
|
|
"yellow": curses.color_pair(2) | curses.A_BOLD,
|
|
"green": curses.color_pair(3) | curses.A_BOLD,
|
|
"red": curses.color_pair(4) | curses.A_BOLD,
|
|
"magenta": curses.color_pair(5) | curses.A_BOLD,
|
|
"head_sel": curses.color_pair(6) | curses.A_BOLD,
|
|
"row_sel": curses.color_pair(7) | curses.A_BOLD,
|
|
"warn": curses.color_pair(8) | curses.A_BOLD,
|
|
"dim": curses.A_DIM,
|
|
}
|
|
return mapping.get(name, 0)
|
|
except Exception:
|
|
return 0
|
|
|
|
def freshness_attr(self, bucket: str) -> int:
|
|
return {"fresh": self._attr("green"),
|
|
"aging": self._attr("yellow"),
|
|
"stale": self._attr("red")}.get(bucket, self._attr("dim"))
|
|
|
|
# -- data -------------------------------------------------------
|
|
|
|
def refresh(self) -> None:
|
|
try:
|
|
self.snapshot = gather_all(timer_names=self.timer_names)
|
|
self.status_msg = "Snapshot %s (auto-refresh %ds; r=refresh)" % (
|
|
datetime.now().strftime("%H:%M:%S"), int(AUTO_REFRESH_S))
|
|
except Exception as e:
|
|
self.snapshot = {}
|
|
self.status_msg = "Refresh failed (showing n/a): %s" % e
|
|
self.last_refresh = time.time()
|
|
self.scroll = 0
|
|
|
|
# -- render helpers ---------------------------------------------
|
|
|
|
def safe_addstr(self, y: int, x: int, text: str, attr: int = 0) -> None:
|
|
h, w = self.stdscr.getmaxyx()
|
|
if 0 <= y < h and 0 <= x < w:
|
|
try:
|
|
self.stdscr.addstr(y, x, text[:max(0, w - x - 1)], attr)
|
|
except Exception:
|
|
pass
|
|
|
|
def _render_header(self, w: int) -> None:
|
|
self.safe_addstr(0, 0, " " * w, self._attr("head_sel"))
|
|
title = " FLEET CONTROL [box fleet-tui] "
|
|
self.safe_addstr(0, 1, title, self._attr("head_sel"))
|
|
self.safe_addstr(1, 0, " " * w, self._attr("dim"))
|
|
col = 1
|
|
for idx, tab_name in enumerate(self.tabs):
|
|
pill = " [%s] " % tab_name
|
|
attr = self._attr("head_sel") if idx == self.current_tab else self._attr("dim")
|
|
self.safe_addstr(1, col, pill, attr)
|
|
col += len(pill) + 1
|
|
self.safe_addstr(2, 0, "-" * w, self._attr("dim"))
|
|
|
|
def _render_footer(self, h: int, w: int) -> None:
|
|
self.safe_addstr(h - 2, 0, "-" * w, self._attr("dim"))
|
|
if self.current_tab == self.RUNTIMES_TAB:
|
|
hints = (" 1-7/Tab: Tabs j/k: Select R: Reconcile "
|
|
"o: Attach ?: Help q: Quit")
|
|
elif self.current_tab == self.WORK_TAB:
|
|
hints = (" 1-7/Tab: Tabs W: Dispatch F: Refocus "
|
|
"L: Sweep H: Heal r: Refresh ?: Help q: Quit")
|
|
else:
|
|
hints = (" 1-7/Tab: Tabs j/k: Scroll r: Refresh "
|
|
"?: Help q: Quit")
|
|
self.safe_addstr(h - 1, 1, self.status_msg[: w - 2], self._attr("dim"))
|
|
if len(self.status_msg) + len(hints) + 2 < w:
|
|
self.safe_addstr(h - 1, w - len(hints) - 1, hints, self._attr("dim"))
|
|
|
|
def _body(self, h: int, w: int, title: str,
|
|
lines: List[Tuple[str, str]]) -> None:
|
|
self.safe_addstr(3, 2, title, self._attr("cyan"))
|
|
self.safe_addstr(4, 2, "-" * (w - 4), self._attr("dim"))
|
|
max_rows = max(0, h - 8)
|
|
visible = lines[self.scroll:self.scroll + max_rows]
|
|
for i, (text, attr_name) in enumerate(visible):
|
|
self.safe_addstr(5 + i, 2, text, self._attr(attr_name))
|
|
if self.scroll > 0:
|
|
self.safe_addstr(5, w - 6, "^more", self._attr("dim"))
|
|
if self.scroll + max_rows < len(lines):
|
|
self.safe_addstr(h - 3, w - 6, "vmore", self._attr("dim"))
|
|
|
|
# -- per-tab renderers ------------------------------------------
|
|
|
|
def _render_timers(self, h: int, w: int) -> None:
|
|
data = self.snapshot.get("timers", {})
|
|
rows = data.get("rows", [])
|
|
lines: List[Tuple[str, str]] = [
|
|
("%-32s %-11s %-9s %-11s %-9s %s"
|
|
% ("TIMER", "NEXT", "IN", "LAST", "AGO", "ACTIVATES"), "dim"),
|
|
]
|
|
for r in rows:
|
|
if not isinstance(r, dict):
|
|
continue
|
|
next_val = _safe_str(r.get("next"), NA)
|
|
last_val = _safe_str(r.get("last"), NA)
|
|
missing = (next_val == NA and last_val == NA)
|
|
attr = "dim" if missing else ("yellow" if next_val == NA else "normal")
|
|
lines.append((
|
|
"%-32s %-11s %-9s %-11s %-9s %s" % (
|
|
_safe_str(r.get("timer"), "?")[:32], next_val[:11],
|
|
_safe_str(r.get("next_rel"), NA)[:9], last_val[:11],
|
|
_safe_str(r.get("last_rel"), NA)[:9],
|
|
_safe_str(r.get("activates"), NA)[:28]), attr))
|
|
if not rows:
|
|
lines.append(("(No fleet timers discovered in systemd/)", "dim"))
|
|
if data.get("note"):
|
|
lines.append(("", "normal"))
|
|
lines.append(("note: %s" % _safe_str(data["note"])[: max(0, w - 10)], "yellow"))
|
|
self._body(h, w, "USER TIMER STATES (NEXT/LAST)", lines)
|
|
|
|
def _render_harvest(self, h: int, w: int) -> None:
|
|
data = self.snapshot.get("harvest", {})
|
|
rows = data.get("rows", [])
|
|
lines: List[Tuple[str, str]] = [
|
|
("%-6s %-20s %-24s %-11s %-9s %s"
|
|
% ("AGENT", "THREAD", "WATERMARK", "LAST", "AGO", "FRESH"), "dim"),
|
|
]
|
|
for r in rows:
|
|
if not isinstance(r, dict):
|
|
continue
|
|
fresh = _safe_str(r.get("freshness"), "unknown")
|
|
if fresh == "fresh":
|
|
attr = "green"
|
|
elif fresh == "aging":
|
|
attr = "yellow"
|
|
elif fresh == "stale":
|
|
attr = "red"
|
|
elif _safe_str(r.get("watermark"), NA) == NA:
|
|
attr = "dim"
|
|
else:
|
|
attr = "normal"
|
|
lines.append((
|
|
"%-6s %-20s %-24s %-11s %-9s %s" % (
|
|
_safe_str(r.get("agent"), "?"), _safe_str(r.get("thread_name"), "?")[:20],
|
|
_safe_str(r.get("watermark"), NA)[:24], _safe_str(r.get("last"), NA)[:11],
|
|
_safe_str(r.get("last_rel"), NA)[:9], fresh), attr))
|
|
if not rows:
|
|
lines.append(("(No harvest watermarks found in %s)"
|
|
% WATERMARKS_FILE.name, "dim"))
|
|
if data.get("note"):
|
|
lines.append(("", "normal"))
|
|
lines.append(("note: %s" % _safe_str(data["note"])[: max(0, w - 10)], "yellow"))
|
|
self._body(h, w, "HARVEST WATERMARKS PER AGENT (freshness highlighted)", lines)
|
|
|
|
def _render_followups(self, h: int, w: int) -> None:
|
|
data = self.snapshot.get("followups", {})
|
|
counts = data.get("counts", {})
|
|
lines: List[Tuple[str, str]] = [
|
|
("Queue depth: %d tracked | overdue: %d" % (
|
|
data.get("total", 0), data.get("overdue", 0)), "cyan"),
|
|
("", "normal"),
|
|
("%-12s %s" % ("STATUS", "COUNT"), "dim"),
|
|
]
|
|
for st in ("pending", "escalated", "resolved", "other"):
|
|
if st not in counts:
|
|
continue
|
|
attr = {"pending": "yellow", "escalated": "red",
|
|
"resolved": "green"}.get(st, "dim")
|
|
lines.append(("%-12s %d" % (st.upper(), counts[st]), attr))
|
|
lines.append(("", "normal"))
|
|
by_rec = data.get("by_recipient", {})
|
|
if by_rec:
|
|
lines.append(("Pending by recipient:", "dim"))
|
|
for agent, n in sorted(by_rec.items()):
|
|
lines.append((" %-8s %d" % (agent, n), "normal"))
|
|
lines.append(("", "normal"))
|
|
oldest = data.get("oldest_pending")
|
|
if isinstance(oldest, dict):
|
|
lines.append(("Oldest pending: %s %s %s sent %s (%s)" % (
|
|
_safe_str(oldest.get("dm_id"), "?"), _safe_str(oldest.get("route"), "?"),
|
|
_safe_str(oldest.get("target"), "?"), _safe_str(oldest.get("sent"), NA),
|
|
_safe_str(oldest.get("sent_rel"), NA)), "yellow"))
|
|
else:
|
|
lines.append(("Oldest pending: none", "dim"))
|
|
if data.get("note"):
|
|
lines.append(("", "normal"))
|
|
lines.append(("note: %s" % _safe_str(data["note"])[: max(0, w - 10)], "yellow"))
|
|
self._body(h, w, "FOLLOW-UP QUEUE COUNTS BY STATUS", lines)
|
|
|
|
def _render_approvals(self, h: int, w: int) -> None:
|
|
data = self.snapshot.get("approvals", {})
|
|
pending = data.get("pending", [])
|
|
lines: List[Tuple[str, str]] = [
|
|
("Pending: %d | checked: %d nodes | unreachable: %s" % (
|
|
data.get("total_pending", 0), data.get("checked", 0),
|
|
", ".join(data.get("unreachable", [])) or "none"), "cyan"),
|
|
("", "normal"),
|
|
]
|
|
if not pending:
|
|
lines.append(("(No pending approvals — or probe returned n/a.)", "dim"))
|
|
else:
|
|
lines.append(("%-6s %-12s %-32s %s"
|
|
% ("NODE", "STATUS", "TARGET", "TITLE"), "dim"))
|
|
for p in pending:
|
|
if not isinstance(p, dict):
|
|
continue
|
|
lines.append((
|
|
"%-6s %-12s %-32s %s" % (
|
|
_safe_str(p.get("node"), "?"), _safe_str(p.get("status"), "?")[:12],
|
|
_safe_str(p.get("target"), "-")[:32],
|
|
_safe_str(p.get("title"), "")[: max(0, w - 58)]), "red"))
|
|
if data.get("note"):
|
|
lines.append(("", "normal"))
|
|
lines.append(("note: %s" % _safe_str(data["note"])[: max(0, w - 10)], "yellow"))
|
|
self._body(h, w, "PENDING APPROVALS", lines)
|
|
|
|
def _render_activity(self, h: int, w: int) -> None:
|
|
data = self.snapshot.get("activity", {})
|
|
rows = data.get("rows", [])
|
|
lines: List[Tuple[str, str]] = [
|
|
("%-6s %-11s %-9s %-6s %-7s %s"
|
|
% ("AGENT", "LAST", "AGO", "SRC", "FRESH", "DETAIL"), "dim"),
|
|
]
|
|
for r in rows:
|
|
if not isinstance(r, dict):
|
|
continue
|
|
fresh = _safe_str(r.get("freshness"), "unknown")
|
|
attr = "green" if fresh == "fresh" else (
|
|
"yellow" if fresh == "aging" else (
|
|
"red" if fresh == "stale" else "dim"))
|
|
lines.append((
|
|
"%-6s %-11s %-9s %-6s %-7s %s" % (
|
|
_safe_str(r.get("agent"), "?"), _safe_str(r.get("last"), NA)[:11],
|
|
_safe_str(r.get("last_rel"), NA)[:9], _safe_str(r.get("source"), NA)[:6],
|
|
fresh, _safe_str(r.get("detail"), NA)[: max(0, w - 48)]), attr))
|
|
if data.get("note"):
|
|
lines.append(("", "normal"))
|
|
lines.append(("note: %s" % _safe_str(data["note"])[: max(0, w - 10)], "yellow"))
|
|
self._body(h, w, "PER-AGENT LAST ACTIVITY", lines)
|
|
|
|
def _render_runtimes(self, h: int, w: int) -> None:
|
|
data = self.snapshot.get("runtimes", {})
|
|
queue = data.get("queue", {})
|
|
plan = data.get("plan", {})
|
|
claimed = queue.get("claimed", {}) or {}
|
|
claimed_txt = (", ".join("%s:%d" % (s, n)
|
|
for s, n in sorted(claimed.items()))
|
|
or "none")
|
|
lines: List[Tuple[str, str]] = [
|
|
("Manifest agents | last reconcile run: %s"
|
|
% data.get("last_run", NA), "cyan"),
|
|
("", "normal"),
|
|
("%-26s %-8s %-5s %-5s %-11s %-10s %-7s %s"
|
|
% ("SESSION", "HAT", "LIVE", "PANE", "STATE", "MODE",
|
|
"WATCHER", "BRIEFED"), "dim"),
|
|
]
|
|
items: List[Dict[str, str]] = []
|
|
line_of_item: Dict[int, int] = {}
|
|
for agent in data.get("agents") or []:
|
|
if not isinstance(agent, dict):
|
|
continue
|
|
live = agent.get("live")
|
|
briefed = agent.get("briefed")
|
|
if isinstance(briefed, bool):
|
|
briefed_txt = "yes" if briefed else "NO"
|
|
else:
|
|
briefed_txt = str(briefed) if briefed is not None else "-"
|
|
attr = "dim" if not agent.get("enabled") else (
|
|
"red" if not live else (
|
|
"yellow" if briefed_txt in ("NO", "stale")
|
|
else "normal"))
|
|
idx = len(items)
|
|
if live:
|
|
items.append({"socket": _safe_str(agent.get("socket"), "?"),
|
|
"session": _safe_str(agent.get("session"), "?")})
|
|
line_of_item[idx] = len(lines)
|
|
else:
|
|
idx = -1
|
|
marker = ">" if idx == self.rt_sel else " "
|
|
lines.append((
|
|
"%s%-25s %-8s %-5s %-5s %-11s %-10s %-7s %s" % (
|
|
marker, _safe_str(agent.get("session"), "?")[:25],
|
|
_safe_str(agent.get("hat"), "?")[:8],
|
|
"YES" if live else "NO",
|
|
_safe_str(agent.get("pane"), NA)[:5],
|
|
_safe_str(agent.get("state"), NA)[:11],
|
|
_safe_str(agent.get("mode"), NA)[:10],
|
|
"ALIVE" if agent.get("watcher") else "-",
|
|
briefed_txt), attr))
|
|
strays = data.get("strays") or []
|
|
if strays:
|
|
lines.append(("", "normal"))
|
|
lines.append(("Stray sessions (fleet sockets, not in manifest):",
|
|
"dim"))
|
|
for stray in strays:
|
|
if not isinstance(stray, dict):
|
|
continue
|
|
idx = len(items)
|
|
items.append({"socket": _safe_str(stray.get("socket"), "?"),
|
|
"session": _safe_str(stray.get("session"), "?")})
|
|
line_of_item[idx] = len(lines)
|
|
marker = ">" if idx == self.rt_sel else " "
|
|
lines.append((
|
|
"%s%-25s %-8s %-5s %-5s %-11s %-10s %-7s %s" % (
|
|
marker, _safe_str(stray.get("session"), "?")[:25],
|
|
"stray",
|
|
"YES", _safe_str(stray.get("pane"), NA)[:5],
|
|
_safe_str(stray.get("state"), NA)[:11], "", "",
|
|
_safe_str(stray.get("cmd"))[:20]), "yellow"))
|
|
lines.append(("", "normal"))
|
|
lines.append(("Task queue: pending %d | claimed %s | done %d"
|
|
% (queue.get("pending", 0), claimed_txt[: max(0, w - 48)],
|
|
queue.get("done", 0)), "cyan"))
|
|
plan_bits = []
|
|
for key in ("launch", "brief", "nudge", "failed"):
|
|
names = plan.get(key, []) or []
|
|
if names:
|
|
plan_bits.append("%s: %s" % (key, ", ".join(str(n) for n in names)))
|
|
lines.append(("", "normal"))
|
|
lines.append(("Reconcile plan (dry-run): %s"
|
|
% ("; ".join(plan_bits) or "steady state"), "cyan"))
|
|
for err in data.get("plan_errors", []) or []:
|
|
lines.append(("plan error: %s" % str(err)[: max(0, w - 14)], "red"))
|
|
if data.get("note"):
|
|
lines.append(("", "normal"))
|
|
lines.append(("note: %s" % _safe_str(data["note"])[: max(0, w - 10)], "yellow"))
|
|
self.rt_items = items
|
|
if self.rt_items:
|
|
self.rt_sel = max(0, min(self.rt_sel, len(self.rt_items) - 1))
|
|
sel_line = line_of_item.get(self.rt_sel)
|
|
if sel_line is not None and sel_line < len(lines):
|
|
text, _ = lines[sel_line]
|
|
lines[sel_line] = (text, "row_sel")
|
|
self._body(h, w, "AGENT SOCKETS/RUNTIMES (j/k select, R reconcile, o attach)", lines)
|
|
|
|
def _action_reconcile(self) -> None:
|
|
self.status_msg = "Running reconciler..."
|
|
self.stdscr.refresh()
|
|
ok, msg = do_reconcile_now()
|
|
self.refresh()
|
|
self.status_msg = msg
|
|
|
|
def _action_attach(self) -> None:
|
|
if not self.rt_items:
|
|
self.status_msg = "Nothing to attach (no live sessions listed)"
|
|
return
|
|
self.rt_sel = max(0, min(self.rt_sel, len(self.rt_items) - 1))
|
|
item = self.rt_items[self.rt_sel]
|
|
argv = attach_argv(item["socket"], item["session"])
|
|
try:
|
|
curses.endwin()
|
|
r = subprocess.run(argv, timeout=86400)
|
|
rc = r.returncode
|
|
except Exception as e:
|
|
rc = -1
|
|
err = str(e)
|
|
try:
|
|
self.stdscr.refresh()
|
|
self.stdscr.nodelay(True)
|
|
self.stdscr.keypad(True)
|
|
try:
|
|
curses.curs_set(0)
|
|
except Exception:
|
|
pass
|
|
except Exception:
|
|
pass
|
|
if rc == 0:
|
|
self.status_msg = "Detached from %s (r to refresh)" % item["session"]
|
|
elif rc == -1:
|
|
self.status_msg = "Attach failed: %s" % err
|
|
else:
|
|
self.status_msg = ("Attach to %s exited rc=%d "
|
|
"(session may be gone)" % (item["session"], rc))
|
|
self.refresh()
|
|
|
|
def _render_work(self, h: int, w: int) -> None:
|
|
data = self.snapshot.get("work", {})
|
|
cog = data.get("cognitive", {})
|
|
issues = data.get("issues", [])
|
|
loop_active = data.get("loopback_active", False)
|
|
|
|
lines: List[Tuple[str, str]] = [
|
|
("COGNITIVE CONSENSUS & WORK PIPELINE | Loopback Daemon: %s"
|
|
% ("ACTIVE" if loop_active else "INACTIVE"), "cyan" if loop_active else "yellow"),
|
|
("", "normal"),
|
|
("%-6s %-16s %-8s %-24s %s"
|
|
% ("AGENT", "COGNITIVE STATE", "LOCK", "ACTIVE URL", "TITLE"), "dim"),
|
|
]
|
|
|
|
for agent in FLEET_AGENTS:
|
|
st = cog.get(agent) or {}
|
|
badge = _safe_str(st.get("badge") or st.get("status"), "UNKNOWN")
|
|
lock = "LOCKED" if st.get("cognitive_lock") else "OPEN"
|
|
scr = st.get("screen") or {}
|
|
url = _safe_str(scr.get("url"), "-")
|
|
title = _safe_str(scr.get("title"), "-")
|
|
lock_attr = "red" if lock == "LOCKED" else "green"
|
|
lines.append((
|
|
"%-6s %-16s %-8s %-24s %s" % (
|
|
agent, badge[:16], lock, url[:24], title[: max(0, w - 60)]
|
|
), lock_attr
|
|
))
|
|
|
|
lines.append(("", "normal"))
|
|
lines.append(("BUILD TICKETS (GITEA REPO: super/box)", "cyan"))
|
|
lines.append(("%-6s %-10s %-10s %-12s %s"
|
|
% ("ISSUE", "STATE", "ASSIGNEE", "CREATED", "TITLE"), "dim"))
|
|
if not issues:
|
|
lines.append(("(No recent build tickets found on tea.muse-dev.online)", "dim"))
|
|
else:
|
|
for idx, iss in enumerate((issues or [])[:8]):
|
|
if not isinstance(iss, dict):
|
|
continue
|
|
inum = "#%s" % _safe_str(iss.get("number"), "?")
|
|
istate = _safe_str(iss.get("state"), "").upper()
|
|
asg = _safe_str((iss.get("assignee") or {}).get("username"), "unassigned")
|
|
created = _safe_str(iss.get("created_at"))[:10]
|
|
title = _safe_str(iss.get("title"), "")
|
|
marker = ">" if idx == self.work_sel else " "
|
|
i_attr = "green" if istate == "CLOSED" else "yellow"
|
|
lines.append((
|
|
"%s%-5s %-10s %-10s %-12s %s" % (
|
|
marker, inum, istate, asg[:10], created, title[: max(0, w - 45)]
|
|
), i_attr
|
|
))
|
|
|
|
self._body(h, w, "WORK PIPELINE & COGNITIVE CONSENSUS (W dispatch, F refocus, L sweep, H heal)", lines)
|
|
|
|
def _action_dispatch_work(self) -> None:
|
|
self.status_msg = "Dispatching work ticket via box work..."
|
|
self.stdscr.refresh()
|
|
cmd = [sys.executable, str(BIN_DIR / "box-work.py"), "start", "--agent", "dev", "Run autonomous build task"]
|
|
rc, out = _run(cmd, timeout=30)
|
|
self.refresh()
|
|
first = out.splitlines()[0] if out else "Done"
|
|
self.status_msg = "Dispatch result: %s" % first[:60]
|
|
|
|
def _action_refocus(self) -> None:
|
|
agent = FLEET_AGENTS[self.work_sel % len(FLEET_AGENTS)]
|
|
self.status_msg = "Refocusing @%s to Main Chat..." % agent.upper()
|
|
self.stdscr.refresh()
|
|
cmd = [sys.executable, str(BIN_DIR / "agent-cognitive-probe.py"), "nav-main", agent]
|
|
rc, out = _run(cmd, timeout=15)
|
|
self.refresh()
|
|
self.status_msg = "Refocus @%s: %s" % (agent.upper(), "OK (🟢 IDLE)" if rc == 0 else "FAILED")
|
|
|
|
def _action_loopback(self) -> None:
|
|
self.status_msg = "Running cognitive loopback sweep..."
|
|
self.stdscr.refresh()
|
|
cmd = [sys.executable, str(BIN_DIR / "box-readback-loopback.py"), "sweep"]
|
|
rc, out = _run(cmd, timeout=30)
|
|
self.refresh()
|
|
self.status_msg = "Loopback sweep: %s" % ("Completed OK" if rc == 0 else "Completed with warnings")
|
|
|
|
def _action_heal(self) -> None:
|
|
agent = FLEET_AGENTS[self.work_sel % len(FLEET_AGENTS)]
|
|
self.status_msg = "Auto-healing cognitive state for @%s..." % agent.upper()
|
|
self.stdscr.refresh()
|
|
cmd = [sys.executable, str(BIN_DIR / "box-work.py"), "heal", agent]
|
|
rc, out = _run(cmd, timeout=20)
|
|
self.refresh()
|
|
self.status_msg = "Heal @%s: %s" % (agent.upper(), "Restored" if rc == 0 else "Alerts")
|
|
|
|
def _render_help(self, h: int, w: int) -> None:
|
|
modal_w = max(10, min(64, w - 6))
|
|
modal_h = max(5, min(15, h - 2))
|
|
top = max(0, (h - modal_h) // 2)
|
|
left = max(0, (w - modal_w) // 2)
|
|
for y in range(top, top + modal_h):
|
|
self.safe_addstr(y, left, " " * modal_w, self._attr("row_sel"))
|
|
self.safe_addstr(top, left, "+" + "-" * (modal_w - 2) + "+",
|
|
self._attr("cyan"))
|
|
for y in range(top + 1, top + modal_h - 1):
|
|
self.safe_addstr(y, left, "|", self._attr("cyan"))
|
|
self.safe_addstr(y, left + modal_w - 1, "|", self._attr("cyan"))
|
|
self.safe_addstr(top + modal_h - 1, left, "+" + "-" * (modal_w - 2) + "+",
|
|
self._attr("cyan"))
|
|
self.safe_addstr(top + 1, left + 3, "FLEET CONTROL HELP",
|
|
self._attr("cyan"))
|
|
for i, hint in enumerate([
|
|
"1-7 / Tab: switch surfaces (1-5 read-only)",
|
|
"j/k / Up/Down: scroll (select on Runtimes & Work)",
|
|
"r: refresh snapshot now",
|
|
"R (Runtimes): run the reconciler now",
|
|
"o / Enter (Runtimes): attach to selection",
|
|
"W (Work): dispatch build ticket",
|
|
"F (Work): refocus agent to Main Chat",
|
|
"L (Work): loopback sweep",
|
|
"H (Work): heal agent",
|
|
"q / Esc: quit (Esc closes help first)",
|
|
"",
|
|
"Missing data shows as 'n/a' — never a crash.",
|
|
]):
|
|
self.safe_addstr(top + 3 + i, left + 4, hint, self._attr("normal"))
|
|
|
|
# -- input + main loop ------------------------------------------
|
|
|
|
def _handle_key(self, ch: int) -> bool:
|
|
if ch in (3, 4): # Ctrl+C / Ctrl+D
|
|
return False
|
|
if self.show_help:
|
|
if ch in (27, ord("q"), ord("Q"), ord("?")):
|
|
self.show_help = False
|
|
return True
|
|
if ch in (ord("q"), ord("Q")):
|
|
return False
|
|
if ch in (27, ord("?")):
|
|
self.show_help = True
|
|
return True
|
|
if ch in (ord("1"), ord("2"), ord("3"), ord("4"), ord("5"),
|
|
ord("6"), ord("7")):
|
|
self.current_tab = ch - ord("1")
|
|
self.scroll = 0
|
|
return True
|
|
if ch == ord("\t"):
|
|
self.current_tab = (self.current_tab + 1) % len(self.tabs)
|
|
self.scroll = 0
|
|
return True
|
|
if ch in (ord("j"), curses.KEY_DOWN):
|
|
if self.current_tab == self.RUNTIMES_TAB and self.rt_items:
|
|
self.rt_sel = min(len(self.rt_items) - 1, self.rt_sel + 1)
|
|
elif self.current_tab == self.WORK_TAB:
|
|
self.work_sel = min(7, self.work_sel + 1)
|
|
else:
|
|
self.scroll += 1
|
|
return True
|
|
if ch in (ord("k"), curses.KEY_UP):
|
|
if self.current_tab == self.RUNTIMES_TAB and self.rt_items:
|
|
self.rt_sel = max(0, self.rt_sel - 1)
|
|
elif self.current_tab == self.WORK_TAB:
|
|
self.work_sel = max(0, self.work_sel - 1)
|
|
else:
|
|
self.scroll = max(0, self.scroll - 1)
|
|
return True
|
|
if ch == ord("r"):
|
|
self.refresh()
|
|
return True
|
|
if ch == ord("R"):
|
|
if self.current_tab == self.RUNTIMES_TAB:
|
|
self._action_reconcile()
|
|
else:
|
|
self.refresh()
|
|
return True
|
|
if ch in (ord("o"), ord("O"), 10, 13, curses.KEY_ENTER):
|
|
if self.current_tab == self.RUNTIMES_TAB:
|
|
self._action_attach()
|
|
return True
|
|
if ch in (ord("w"), ord("W")):
|
|
if self.current_tab == self.WORK_TAB:
|
|
self._action_dispatch_work()
|
|
return True
|
|
if ch in (ord("f"), ord("F")):
|
|
if self.current_tab == self.WORK_TAB:
|
|
self._action_refocus()
|
|
return True
|
|
if ch in (ord("l"), ord("L")):
|
|
if self.current_tab == self.WORK_TAB:
|
|
self._action_loopback()
|
|
return True
|
|
if ch in (ord("h"), ord("H")):
|
|
if self.current_tab == self.WORK_TAB:
|
|
self._action_heal()
|
|
return True
|
|
return True
|
|
|
|
def run(self) -> None:
|
|
renderers = [self._render_timers, self._render_harvest,
|
|
self._render_followups, self._render_approvals,
|
|
self._render_activity, self._render_runtimes,
|
|
self._render_work]
|
|
while True:
|
|
h, w = self.stdscr.getmaxyx()
|
|
self.stdscr.erase()
|
|
self._render_header(w)
|
|
try:
|
|
renderers[self.current_tab](h, w)
|
|
except Exception as e:
|
|
self.safe_addstr(5, 4, "Render error (n/a): %s" % e,
|
|
self._attr("red"))
|
|
self._render_footer(h, w)
|
|
if self.show_help:
|
|
self._render_help(h, w)
|
|
self.stdscr.refresh()
|
|
try:
|
|
ch = self.stdscr.getch()
|
|
if ch != -1 and not self._handle_key(ch):
|
|
break
|
|
except KeyboardInterrupt:
|
|
break
|
|
if time.time() - self.last_refresh > AUTO_REFRESH_S:
|
|
self.refresh()
|
|
time.sleep(0.05)
|
|
|
|
|
|
def main(argv: Optional[List[str]] = None) -> int:
|
|
argv = list(sys.argv[1:] if argv is None else argv)
|
|
if "--once" in argv:
|
|
snap = gather_all()
|
|
if "--json" in argv:
|
|
print(json.dumps(snap, indent=2, default=str))
|
|
else:
|
|
for surface, data in snap.items():
|
|
print("== %s ==" % surface.upper())
|
|
print(json.dumps(data, indent=2, default=str)[:2000])
|
|
return 0
|
|
curses.wrapper(lambda stdscr: BoxFleetTUI(stdscr).run())
|
|
return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
sys.exit(main())
|