From 3ddd9ca08c4dc65f41bb075ba6db0d96a4ed8da9 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 14:13:36 -0400 Subject: [PATCH] feat(agent-team): read-only LAN status dashboard module (status_page.py) --- agent-team/agent_team/status_page.py | 684 +++++++++++++++++++++++++++ 1 file changed, 684 insertions(+) create mode 100644 agent-team/agent_team/status_page.py diff --git a/agent-team/agent_team/status_page.py b/agent-team/agent_team/status_page.py new file mode 100644 index 0000000..55c5c3f --- /dev/null +++ b/agent-team/agent_team/status_page.py @@ -0,0 +1,684 @@ +"""LAN-only, READ-ONLY web status dashboard for the agent-team coordinator. + +The agent-team is a single durable LangGraph pipeline; "agents" here are the +per-stage model roles, and what an operator actually wants to see is *which task +is in which phase* and *which tasks are blocked on the human gate*. This module +serves that view as a tiny self-refreshing HTML page so Adam can glance at the +queue without SSHing to the R720 and running ``run-team.py list``. + +Posture (read this before changing anything): + +* **READ-ONLY.** The page exposes no mutating endpoints. The SQLite ledger is + opened read-only (``mode=ro`` / ``immutable``-safe URI) and never written. It + reads the same two durable sources the coordinator owns: + - ``pending_questions`` (the human-interaction lifecycle ledger) for "who is + waiting on what" — reused via :mod:`agent_team.ledger`. + - the LangGraph SQLite checkpoint (``channel_values`` holds ``task``, + ``current_phase``, ``status``) for each task's phase/description — read + through the same :class:`SqliteSaver` the coordinator uses in + :func:`agent_team.graph.build_sqlite_checkpointer`. +* **Sensitive content.** Task descriptions are operator input and may contain + sensitive detail (repo names, ticket content, remediation context). The page + renders them in plain text with HTML-escaping but applies NO authentication. + It must therefore stay **LAN/VPN-only and never public**. The sh-secrev R720 + VM (10.10.60.120, VLAN 60) has no public NIC and sits behind the UniFi + firewall, so binding ``0.0.0.0`` exposes it to the LAN/VPN only. No secrets, + tokens, channel refs, or answer bodies are rendered. +* **Fail-safe.** A missing, locked, or unreadable DB renders a friendly + "no data" page; the server never crashes on a bad read. + +Why stdlib ``http.server`` and not FastAPI (which is already a dep): the page is +a single read-only GET that returns server-rendered HTML on a 10s meta-refresh. +FastAPI/uvicorn buys nothing here (no JSON contract, no auth, no async I/O) and +would couple this LAN dashboard to the optional WS1 HTTP-API dependency set. +stdlib ``http.server`` keeps it zero-new-deps and importable anywhere the +``agent_team`` package is, which also keeps :func:`render_html` trivially +unit-testable without binding a socket. + +Config (all env, with defaults): + +* ``AGENT_TEAM_DB`` — path to the SQLite ledger + (default ``agent-team/state/agent_team.sqlite`` relative to this package). +* ``AGENT_TEAM_STATUS_HOST`` — bind host (default ``0.0.0.0``; LAN/VPN-only by + network posture, see above). +* ``AGENT_TEAM_STATUS_PORT`` — bind port (default ``8770``). +""" + +from __future__ import annotations + +import html +import os +import sqlite3 +from collections import Counter +from dataclasses import dataclass, field +from datetime import datetime, timezone +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path +from typing import Any + +# The package root is .../agent-team/agent_team; the default ledger lives at +# .../agent-team/state/agent_team.sqlite (mirrors run-team.py's _DEFAULT_DB). +_PACKAGE_DIR = Path(__file__).resolve().parent +_DEFAULT_DB = _PACKAGE_DIR.parent / "state" / "agent_team.sqlite" + +_DEFAULT_HOST = "0.0.0.0" +_DEFAULT_PORT = 8770 + +# How long the browser waits before re-requesting (seconds). +_REFRESH_SECONDS = 10 + +# Phases we treat as "still progressing" for the queue summary. PARKED/DONE/ +# FAILED are terminal-ish; everything else is in flight. +_ACTIVE_STATUSES = {"active", "waiting_human"} +_PARKED_STATUSES = {"parked"} + + +@dataclass(frozen=True) +class TaskView: + """A read-only, render-ready view of one pipeline task. + + Built from the LangGraph checkpoint's ``channel_values`` joined to any open + pending question for that thread. ``short_id`` is the first 8 chars of the + thread_id for compact display; ``thread_id`` is retained for completeness. + """ + + thread_id: str + short_id: str + task: str + current_phase: str + status: str + waiting: bool = False + waiting_since: str | None = None + + +@dataclass(frozen=True) +class Snapshot: + """Everything :func:`render_html` needs, as plain data (no DB handle). + + Keeping this pure and serializable is what makes the renderer unit-testable + without a socket or a live database. + """ + + generated_at: str + db_path: str + ok: bool = True + error: str | None = None + tasks: list[TaskView] = field(default_factory=list) + question_counts: dict[str, int] = field(default_factory=dict) + recent_spend: list[dict[str, Any]] = field(default_factory=list) + spend_total_usd: float | None = None + budget_available: bool = False + + @property + def active_count(self) -> int: + return sum(1 for t in self.tasks if t.status in _ACTIVE_STATUSES) + + @property + def parked_count(self) -> int: + return sum(1 for t in self.tasks if t.status in _PARKED_STATUSES) + + @property + def waiting_tasks(self) -> list[TaskView]: + return [t for t in self.tasks if t.waiting] + + @property + def progressing_tasks(self) -> list[TaskView]: + return [t for t in self.tasks if not t.waiting] + + +def _utc_now_iso() -> str: + """Current UTC time as a stable, human-readable string.""" + return datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S UTC") + + +def _short(thread_id: str) -> str: + """First 8 chars of a thread_id, for compact display.""" + return thread_id[:8] if thread_id else "(none)" + + +def _ro_connect(db_path: Path) -> sqlite3.Connection: + """Open ``db_path`` strictly READ-ONLY. + + Uses the SQLite ``file:...?mode=ro`` URI so the connection can never create + or write the database. ``sqlite3.Row`` row factory matches the reuse readers + (:mod:`agent_team.ledger`). Raises if the file does not exist (``mode=ro`` + will not create it) — the caller turns that into the friendly no-data page. + """ + uri = f"file:{db_path}?mode=ro" + conn = sqlite3.connect(uri, uri=True, timeout=2.0) + conn.row_factory = sqlite3.Row + return conn + + +def _enumerate_threads(conn: sqlite3.Connection) -> list[str]: + """Return the distinct thread_ids the checkpointer has persisted. + + The LangGraph ``SqliteSaver`` owns the ``checkpoints`` table; we read it + READ-ONLY. If the table does not exist yet (no task ever ran), returns ``[]`` + rather than raising. + """ + try: + rows = conn.execute( + "SELECT DISTINCT thread_id FROM checkpoints ORDER BY thread_id" + ).fetchall() + except sqlite3.OperationalError: + # No checkpoints table yet (fresh DB) — not an error for a status view. + return [] + return [r["thread_id"] for r in rows if r["thread_id"]] + + +def _channel_values_for(saver: Any, thread_id: str) -> dict[str, Any]: + """Read the latest checkpoint's ``channel_values`` for one thread. + + Mirrors the access pattern the coordinator/graph rely on: + ``saver.get({"configurable": {"thread_id": tid}})`` returns the latest + checkpoint, whose ``channel_values`` carries ``task`` / ``current_phase`` / + ``status`` / ``qa_history``. ``.get`` may return a checkpoint mapping + directly or ``None`` if the thread has no committed checkpoint; we tolerate + both and never raise into the page. + """ + try: + checkpoint = saver.get({"configurable": {"thread_id": thread_id}}) + except Exception: + return {} + if not checkpoint: + return {} + values = checkpoint.get("channel_values") if isinstance(checkpoint, dict) else None + return values if isinstance(values, dict) else {} + + +def build_snapshot(db_path: Path | str | None = None) -> Snapshot: + """Read the ledger READ-ONLY and assemble a render-ready :class:`Snapshot`. + + Fail-safe by contract: any DB problem (missing file, locked, missing tables, + optional langgraph package absent) yields ``ok=False`` with a short message + rather than raising. The renderer turns that into the friendly no-data page. + """ + resolved = Path(db_path) if db_path is not None else _resolve_db_path() + generated_at = _utc_now_iso() + + if not resolved.exists(): + return Snapshot( + generated_at=generated_at, + db_path=str(resolved), + ok=False, + error=f"ledger not found at {resolved} (no tasks have run yet?)", + ) + + conn: sqlite3.Connection | None = None + try: + conn = _ro_connect(resolved) + question_counts = _read_question_counts(conn) + open_by_thread = _read_open_questions(conn) + thread_ids = _enumerate_threads(conn) + tasks = _read_tasks(resolved, thread_ids, open_by_thread) + spend, spend_total, budget_available = _read_recent_spend(conn) + except sqlite3.OperationalError as exc: + # Locked / busy / unreadable — never crash, just say so. + return Snapshot( + generated_at=generated_at, + db_path=str(resolved), + ok=False, + error=f"could not read ledger (busy or unreadable): {exc}", + ) + except Exception as exc: # defensive: any unexpected read failure + return Snapshot( + generated_at=generated_at, + db_path=str(resolved), + ok=False, + error=f"unexpected error reading ledger: {exc}", + ) + finally: + if conn is not None: + conn.close() + + return Snapshot( + generated_at=generated_at, + db_path=str(resolved), + ok=True, + tasks=tasks, + question_counts=question_counts, + recent_spend=spend, + spend_total_usd=spend_total, + budget_available=budget_available, + ) + + +def _resolve_db_path() -> Path: + """Resolve the ledger path from ``AGENT_TEAM_DB`` or the package default.""" + env = os.environ.get("AGENT_TEAM_DB") + return Path(env) if env else _DEFAULT_DB + + +def _read_question_counts(conn: sqlite3.Connection) -> dict[str, int]: + """Count pending_questions rows by status (open/answered/expired/...).""" + counts: Counter[str] = Counter() + try: + rows = conn.execute( + "SELECT status, COUNT(*) AS n FROM pending_questions GROUP BY status" + ).fetchall() + except sqlite3.OperationalError: + return {} + for r in rows: + counts[r["status"]] = int(r["n"]) + return dict(counts) + + +def _read_open_questions(conn: sqlite3.Connection) -> dict[str, str]: + """Map thread_id -> earliest ``posted_at`` of an OPEN question on it. + + An open question means that thread is blocked on the human gate. We keep the + oldest posted_at so the page can show how long it has been waiting. + """ + waiting: dict[str, str] = {} + try: + rows = conn.execute( + "SELECT thread_id, posted_at FROM pending_questions " + "WHERE status='open' ORDER BY posted_at ASC" + ).fetchall() + except sqlite3.OperationalError: + return {} + for r in rows: + tid = r["thread_id"] + if tid and tid not in waiting: + waiting[tid] = r["posted_at"] or "" + return waiting + + +def _read_tasks( + db_path: Path, + thread_ids: list[str], + open_by_thread: dict[str, str], +) -> list[TaskView]: + """Build a :class:`TaskView` per thread from the LangGraph checkpoint. + + Opens its own READ-ONLY :class:`SqliteSaver` over the same DB file (the saver + needs its own connection / serde). If the optional langgraph checkpoint + package is absent, returns ``[]`` so the page degrades to ledger-only data + instead of crashing. + """ + if not thread_ids: + return [] + + saver_cm = _readonly_saver(db_path) + if saver_cm is None: + return [] + + tasks: list[TaskView] = [] + try: + with saver_cm as saver: + for tid in thread_ids: + values = _channel_values_for(saver, tid) + status = str(values.get("status") or "unknown") + tasks.append( + TaskView( + thread_id=tid, + short_id=_short(tid), + task=str(values.get("task") or "(no description)"), + current_phase=str(values.get("current_phase") or "unknown"), + status=status, + waiting=tid in open_by_thread, + waiting_since=open_by_thread.get(tid) or None, + ) + ) + except Exception: + # Saver construction/read failed entirely — degrade to ledger-only. + return tasks + # Waiting tasks first, then by phase, for a sensible default ordering. + tasks.sort(key=lambda t: (not t.waiting, t.current_phase, t.short_id)) + return tasks + + +def _readonly_saver(db_path: Path) -> Any: + """Return an *un-entered* SqliteSaver context manager over a RO connection. + + Reuses the same :class:`SqliteSaver` + serde the coordinator builds in + :func:`agent_team.graph.build_checkpoint_serde`, but over a ``mode=ro`` + connection so this dashboard can never write the checkpoint tables. Returns + ``None`` if the optional package is not installed (pre-deploy scaffolding / + test envs without the checkpoint dep), letting the caller degrade to + ledger-only data. + """ + try: + from contextlib import closing, contextmanager + + from langgraph.checkpoint.sqlite import SqliteSaver + + from agent_team.graph import build_checkpoint_serde + except Exception: + return None + + serde = build_checkpoint_serde() + + @contextmanager + def _cm() -> Any: + uri = f"file:{db_path}?mode=ro" + with closing( + sqlite3.connect(uri, uri=True, check_same_thread=False, timeout=2.0) + ) as conn: + yield SqliteSaver(conn, serde=serde) + + return _cm() + + +def _read_recent_spend( + conn: sqlite3.Connection, +) -> tuple[list[dict[str, Any]], float | None, bool]: + """Read recent budget_ledger spend, if the table is present/readable. + + Returns ``(recent_rows, total_usd, available)``. ``available`` is ``False`` + when the table is absent (older DB) — the page then simply omits the budget + panel rather than showing a misleading zero. + """ + try: + rows = conn.execute( + "SELECT day_bucket, stage, model, usd_cost, recorded_at " + "FROM budget_ledger ORDER BY recorded_at DESC LIMIT 10" + ).fetchall() + total_row = conn.execute( + "SELECT COALESCE(SUM(usd_cost), 0.0) AS total FROM budget_ledger" + ).fetchone() + except sqlite3.OperationalError: + return [], None, False + + recent = [ + { + "day_bucket": r["day_bucket"], + "stage": r["stage"], + "model": r["model"], + "usd_cost": float(r["usd_cost"] or 0.0), + "recorded_at": r["recorded_at"], + } + for r in rows + ] + total = float(total_row["total"]) if total_row is not None else 0.0 + return recent, total, True + + +# --- HTML rendering (PURE: snapshot -> str, no I/O). ------------------------- + +_CSS = """ +:root { color-scheme: light dark; } +* { box-sizing: border-box; } +body { + font-family: -apple-system, BlinkMacSystemFont, "Segoe UI", Roboto, Helvetica, + Arial, sans-serif; + margin: 0; padding: 1.5rem; line-height: 1.45; + background: #0f1115; color: #e6e6e6; +} +h1 { font-size: 1.4rem; margin: 0 0 .25rem; } +h2 { font-size: 1.05rem; margin: 1.5rem 0 .5rem; color: #9fb4d6; } +.meta { color: #8a93a6; font-size: .85rem; margin-bottom: 1rem; } +.cards { display: flex; flex-wrap: wrap; gap: .75rem; margin: .5rem 0 1rem; } +.card { + background: #1a1e26; border: 1px solid #2a303c; border-radius: 8px; + padding: .75rem 1rem; min-width: 7rem; +} +.card .num { font-size: 1.6rem; font-weight: 600; } +.card .lbl { color: #8a93a6; font-size: .8rem; text-transform: uppercase; + letter-spacing: .04em; } +table { border-collapse: collapse; width: 100%; margin: .25rem 0 1rem; + font-size: .9rem; } +th, td { text-align: left; padding: .45rem .6rem; border-bottom: 1px solid #262b35; + vertical-align: top; } +th { color: #9fb4d6; font-weight: 600; } +td.mono, .mono { font-family: ui-monospace, SFMono-Regular, Menlo, monospace; } +.badge { display: inline-block; padding: .1rem .5rem; border-radius: 999px; + font-size: .78rem; font-weight: 600; } +.badge.active { background: #163a2b; color: #6ee7a8; } +.badge.waiting_human, .badge.waiting { background: #3a3216; color: #f5d76e; } +.badge.parked { background: #3a1f16; color: #f3a36e; } +.badge.done { background: #1c2b3a; color: #7fb3f5; } +.badge.failed { background: #3a1620; color: #f57f9c; } +.badge.unknown { background: #2a303c; color: #b8c0cf; } +.desc { max-width: 48rem; } +.empty { color: #8a93a6; font-style: italic; padding: .5rem 0; } +.warn { background: #3a1f16; border: 1px solid #5a3322; color: #f3c79e; + padding: 1rem; border-radius: 8px; } +footer { color: #6b7385; font-size: .78rem; margin-top: 2rem; } +""" + + +def _esc(value: object) -> str: + """HTML-escape any value (task descriptions are untrusted operator input).""" + return html.escape(str(value), quote=True) + + +def _badge(text: str, cls: str | None = None) -> str: + css = (cls or text or "unknown").strip().lower().replace(" ", "_") + return f'{_esc(text)}' + + +def _render_task_rows(tasks: list[TaskView]) -> str: + if not tasks: + return 'No tasks.' + out: list[str] = [] + for t in tasks: + out.append( + "" + f'{_esc(t.short_id)}' + f'{_esc(t.task)}' + f"{_badge(t.current_phase)}" + f"{_badge(t.status)}" + "" + ) + return "".join(out) + + +def _render_waiting(tasks: list[TaskView]) -> str: + waiting = [t for t in tasks if t.waiting] + if not waiting: + return '

No tasks are waiting on the human gate.

' + rows: list[str] = [] + for t in waiting: + since = _esc(t.waiting_since) if t.waiting_since else "—" + rows.append( + "" + f'{_esc(t.short_id)}' + f'{_esc(t.task)}' + f"{_badge(t.current_phase)}" + f'{since}' + "" + ) + return ( + "" + "" + + "".join(rows) + + "
ThreadTaskPhaseOpen since (UTC)
" + ) + + +def _render_budget(snapshot: Snapshot) -> str: + if not snapshot.budget_available: + return "" + total = ( + f"${snapshot.spend_total_usd:,.4f}" + if snapshot.spend_total_usd is not None + else "—" + ) + rows: list[str] = [] + for entry in snapshot.recent_spend: + rows.append( + "" + f'{_esc(entry.get("recorded_at"))}' + f"{_esc(entry.get('stage') or '—')}" + f'{_esc(entry.get("model"))}' + f'${float(entry.get("usd_cost") or 0.0):,.4f}' + "" + ) + body = ( + "".join(rows) + if rows + else 'No spend recorded.' + ) + return ( + f"

Budget — total recorded spend {total}

" + "" + "" + body + "
Recorded (UTC)StageModelUSD
" + ) + + +def render_html(snapshot: Snapshot) -> str: + """Render a :class:`Snapshot` to a complete, self-contained HTML document. + + Pure function: no I/O, no DB, no socket — unit-testable against a fake + snapshot. Inline CSS only and a ```` so the page + works fully offline (no external CDN) and auto-refreshes. + """ + head = ( + '' + '' + f'' + "agent-team status" + f"" + ) + header = ( + "

Sea Haven agent-team — coordinator status

" + f'
Last refresh {_esc(snapshot.generated_at)} ' + f"· auto-refresh {_REFRESH_SECONDS}s · " + f'{_esc(snapshot.db_path)} · ' + "read-only, LAN/VPN-only
" + ) + + if not snapshot.ok: + body = ( + '
No data.
' + f"{_esc(snapshot.error or 'ledger unavailable')}
" + "The dashboard is read-only and waits for the coordinator to " + "produce state.
" + ) + return head + header + body + _footer() + "" + + qc = snapshot.question_counts + cards = ( + '
' + f'
{snapshot.active_count}
' + '
Active
' + f'
{snapshot.parked_count}
' + '
Parked
' + f'
{len(snapshot.waiting_tasks)}
' + '
Waiting (gate)
' + f'
{qc.get("open", 0)}
' + '
Open questions
' + f'
{len(snapshot.tasks)}
' + '
Total tasks
' + "
" + ) + + waiting_section = "

Waiting on the human gate

" + _render_waiting( + snapshot.tasks + ) + + tasks_section = ( + "

All tasks

" + "" + "" + + _render_task_rows(snapshot.tasks) + + "
ThreadTaskPhaseStatus
" + ) + + budget_section = _render_budget(snapshot) + + return ( + head + + header + + cards + + waiting_section + + tasks_section + + budget_section + + _footer() + + "" + ) + + +def _footer() -> str: + return ( + "" + ) + + +# --- HTTP server (thin shell around the pure renderer). ---------------------- + + +def _make_handler(db_path: Path) -> type[BaseHTTPRequestHandler]: + """Build a request handler bound to ``db_path`` (closure, no globals).""" + + class _StatusHandler(BaseHTTPRequestHandler): + server_version = "AgentTeamStatus/1.0" + + def do_GET(self) -> None: # noqa: N802 (BaseHTTPRequestHandler API) + # Read-only: every path serves the dashboard. Favicon gets a 204 so + # browsers stop logging 404s; there are no other routes. + if self.path == "/favicon.ico": + self.send_response(204) + self.end_headers() + return + try: + snapshot = build_snapshot(db_path) + page = render_html(snapshot) + except Exception as exc: # last-ditch: never 500 the dashboard + page = render_html( + Snapshot( + generated_at=_utc_now_iso(), + db_path=str(db_path), + ok=False, + error=f"render error: {exc}", + ) + ) + body = page.encode("utf-8") + self.send_response(200) + self.send_header("Content-Type", "text/html; charset=utf-8") + self.send_header("Content-Length", str(len(body))) + # Defense-in-depth headers for an unauthenticated LAN page. + self.send_header("X-Content-Type-Options", "nosniff") + self.send_header("Referrer-Policy", "no-referrer") + self.send_header("Cache-Control", "no-store") + self.end_headers() + self.wfile.write(body) + + def log_message(self, fmt: str, *args: Any) -> None: + # Quiet by default; access logs are noise for a refresh-loop page. + return + + return _StatusHandler + + +def serve( + host: str | None = None, + port: int | None = None, + db_path: Path | str | None = None, +) -> None: + """Start the LAN-only read-only status server (blocking). + + Resolves host/port/db from args or env (``AGENT_TEAM_STATUS_HOST`` default + ``0.0.0.0``, ``AGENT_TEAM_STATUS_PORT`` default ``8770``, ``AGENT_TEAM_DB`` + default the package ledger). Binds and serves until interrupted. The server + is read-only and unauthenticated — see the module docstring on network + posture (LAN/VPN-only, never public). + """ + resolved_host = host or os.environ.get("AGENT_TEAM_STATUS_HOST") or _DEFAULT_HOST + if port is not None: + resolved_port = int(port) + else: + resolved_port = int(os.environ.get("AGENT_TEAM_STATUS_PORT", _DEFAULT_PORT)) + resolved_db = Path(db_path) if db_path is not None else _resolve_db_path() + + handler = _make_handler(resolved_db) + httpd = ThreadingHTTPServer((resolved_host, resolved_port), handler) + print( + f"[status_page] serving READ-ONLY dashboard on " + f"http://{resolved_host}:{resolved_port} (db={resolved_db}); " + "LAN/VPN-only, unauthenticated." + ) + try: + httpd.serve_forever() + except KeyboardInterrupt: + pass + finally: + httpd.server_close() + + +if __name__ == "__main__": # pragma: no cover - manual / systemd entrypoint + serve()