feat(agent-team): LangGraph-introspected topology + read-only dashboard API
topology.py derives the pipeline map (nodes/edges/trees) from the compiled
LangGraph via get_graph() + a NODE_META display sidecar, so new agent nodes
appear automatically and group into trees branching off intake. dashboard.py is a
new read-only FastAPI app (0.0.0.0:8770) serving /api/state (contract preserved +
per-node live state), /api/topology, and /api/task/{id} (validated, timeline +
cost join + partial fallback) plus the built SPA — kept SEPARATE from the authed
api.py. status_page.py is trimmed to the /api/state data layer; the inline
HTML/SVG renderer + stdlib server are retired.
This commit is contained in:
parent
4656ca64b6
commit
322a1f6922
3 changed files with 556 additions and 723 deletions
292
agent-team/agent_team/dashboard.py
Normal file
292
agent-team/agent_team/dashboard.py
Normal file
|
|
@ -0,0 +1,292 @@
|
|||
"""Read-only LAN status dashboard — FastAPI app + JSON API for the WebUI SPA.
|
||||
|
||||
This is the backend for the React/Vite single-page dashboard. It is a SEPARATE
|
||||
surface from :mod:`agent_team.api` (the authenticated, write-capable, 127.0.0.1
|
||||
WS1 API): this app is **read-only**, **unauthenticated**, and bound LAN-wide on
|
||||
:8770 — the same posture as the retired stdlib ``status_page.serve``. Keeping the
|
||||
two apps separate preserves that security boundary (no auth/write surface leaks
|
||||
onto the LAN page, no read-only LAN exposure on the authed app).
|
||||
|
||||
Endpoints (all read-only):
|
||||
|
||||
``GET /api/state``
|
||||
Live overview — the existing :func:`status_page.snapshot_to_dict` payload
|
||||
(contract preserved) PLUS a ``nodes`` map of per-topology-node live state.
|
||||
``GET /api/topology``
|
||||
The pipeline graph (nodes + edges + trees), introspected from the real
|
||||
LangGraph (:func:`agent_team.topology.build_topology`).
|
||||
``GET /api/task/{thread_id}``
|
||||
One task's history: the ``task_transitions`` timeline + per-stage cost
|
||||
(budget_ledger) + Q&A + review verdicts + plan summary from the checkpoint.
|
||||
|
||||
Static SPA assets (``web/dist``) are mounted at ``/`` LAST so the ``/api`` routes
|
||||
take precedence; when no build is present (dev/test) the mount is skipped.
|
||||
|
||||
All DB access is via a strictly READ-ONLY connection — this app can never write
|
||||
the ledger or the checkpoint tables.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import re
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from agent_team import status_page
|
||||
from agent_team.db.transitions import read_transitions
|
||||
from agent_team.topology import build_topology, node_for_phase
|
||||
|
||||
__all__ = ["make_dashboard_app", "serve", "task_detail"]
|
||||
|
||||
_DEFAULT_HOST = "0.0.0.0" # LAN-only read-only dashboard (matches prior :8770)
|
||||
_DEFAULT_PORT = 8770
|
||||
|
||||
# thread_ids are uuid4 hex (32 chars); allow a generous but bounded charset so a
|
||||
# path param can never carry traversal / injection payloads (plan FIX-2). SQL is
|
||||
# parameterized regardless; this rejects junk early with a 400.
|
||||
_THREAD_ID_RE = re.compile(r"^[A-Za-z0-9_-]{1,64}$")
|
||||
|
||||
# Active / parked status groupings, mirroring status_page's stage_state logic so
|
||||
# the per-node live state on /api/state matches the task badges.
|
||||
_ACTIVE_STATUSES = frozenset({"active"})
|
||||
_WAITING_STATUSES = frozenset({"waiting_human"})
|
||||
_PARKED_STATUSES = frozenset({"parked", "failed"})
|
||||
|
||||
# Map a topology node id back to the budget_ledger ``stage`` key for cost join.
|
||||
# Most node ids equal their stage; the P3 build/verify vertices differ.
|
||||
_NODE_TO_STAGE = {"build_node": "build", "verify_node": "verify"}
|
||||
|
||||
|
||||
def _resolve_dist_dir() -> Path:
|
||||
"""Resolve the built SPA directory (``web/dist``), env-overridable."""
|
||||
env = os.environ.get("AGENT_TEAM_WEB_DIST")
|
||||
if env:
|
||||
return Path(env)
|
||||
# dashboard.py -> agent_team/ -> agent-team/ ; the SPA builds to web/dist.
|
||||
return Path(__file__).resolve().parent.parent / "web" / "dist"
|
||||
|
||||
|
||||
def _node_states(snapshot: status_page.Snapshot) -> dict[str, dict[str, Any]]:
|
||||
"""Per-topology-node live state + count, bucketed from the snapshot's tasks.
|
||||
|
||||
``awaiting_human`` wins over ``active`` over ``parked`` so a blocked node
|
||||
reads as blocked at a glance (same precedence as status_page.stage_state).
|
||||
Nodes with no live task are omitted (the SPA defaults them to idle).
|
||||
"""
|
||||
buckets: dict[str, list[Any]] = {}
|
||||
for task in snapshot.tasks:
|
||||
node_id = node_for_phase(task.current_phase)
|
||||
if node_id is None:
|
||||
continue
|
||||
buckets.setdefault(node_id, []).append(task)
|
||||
|
||||
out: dict[str, dict[str, Any]] = {}
|
||||
for node_id, members in buckets.items():
|
||||
if any(t.waiting for t in members):
|
||||
state = "awaiting_human"
|
||||
elif any(t.status in _WAITING_STATUSES for t in members):
|
||||
state = "awaiting_human"
|
||||
elif any(t.status in _PARKED_STATUSES for t in members):
|
||||
state = "parked"
|
||||
elif any(t.status in _ACTIVE_STATUSES for t in members):
|
||||
state = "active"
|
||||
else:
|
||||
state = "idle"
|
||||
out[node_id] = {"state": state, "count": len(members)}
|
||||
return out
|
||||
|
||||
|
||||
def _duration_seconds(entered_at: str | None, exited_at: str | None) -> float | None:
|
||||
"""Seconds between two ISO-8601 stamps, or ``None`` if still open/unparseable."""
|
||||
if not entered_at or not exited_at:
|
||||
return None
|
||||
try:
|
||||
start = datetime.fromisoformat(entered_at)
|
||||
end = datetime.fromisoformat(exited_at)
|
||||
except ValueError:
|
||||
return None
|
||||
return max(0.0, (end - start).total_seconds())
|
||||
|
||||
|
||||
def _cost_by_stage(conn: Any, thread_id: str) -> dict[str, float]:
|
||||
"""Sum budget_ledger usd_cost per stage for one thread (empty if absent)."""
|
||||
try:
|
||||
rows = conn.execute(
|
||||
"SELECT stage, COALESCE(SUM(usd_cost), 0.0) AS total "
|
||||
"FROM budget_ledger WHERE thread_id = ? GROUP BY stage",
|
||||
(thread_id,),
|
||||
).fetchall()
|
||||
except Exception:
|
||||
return {}
|
||||
return {str(r["stage"]): float(r["total"] or 0.0) for r in rows if r["stage"]}
|
||||
|
||||
|
||||
def task_detail(db_path: Path | str, thread_id: str) -> dict[str, Any]:
|
||||
"""Assemble the per-task drill-down payload (read-only, fail-safe).
|
||||
|
||||
Returns the transition timeline joined to per-stage cost, plus the Q&A
|
||||
history / review verdicts / plan summary read from the LangGraph checkpoint.
|
||||
When no ``task_transitions`` rows exist for the thread (a task that predates
|
||||
the v3 ledger), ``partial`` is ``True`` and a best-effort single-step
|
||||
timeline is derived from the checkpoint's ``current_phase`` (plan FIX-3).
|
||||
"""
|
||||
resolved = Path(db_path)
|
||||
if not resolved.exists():
|
||||
return {"ok": False, "error": "ledger not found", "thread_id": thread_id}
|
||||
|
||||
conn = None
|
||||
try:
|
||||
conn = status_page._ro_connect(resolved)
|
||||
rows = read_transitions(conn, thread_id)
|
||||
cost_by_stage = _cost_by_stage(conn, thread_id)
|
||||
except Exception as exc: # never crash the endpoint
|
||||
return {"ok": False, "error": f"read failed: {exc}", "thread_id": thread_id}
|
||||
finally:
|
||||
if conn is not None:
|
||||
conn.close()
|
||||
|
||||
# Checkpoint-derived fields (task text, status, phase, Q&A, plan, verdicts).
|
||||
values: dict[str, Any] = {}
|
||||
saver_cm = status_page._readonly_saver(resolved)
|
||||
if saver_cm is not None:
|
||||
try:
|
||||
with saver_cm as saver:
|
||||
values = status_page._channel_values_for(saver, thread_id)
|
||||
except Exception:
|
||||
values = {}
|
||||
|
||||
def _stage_cost(node_id: str) -> float:
|
||||
return cost_by_stage.get(_NODE_TO_STAGE.get(node_id, node_id), 0.0)
|
||||
|
||||
partial = not rows
|
||||
timeline: list[dict[str, Any]] = []
|
||||
if rows:
|
||||
for r in rows:
|
||||
timeline.append(
|
||||
{
|
||||
"from_phase": r["from_phase"],
|
||||
"to_phase": r["to_phase"],
|
||||
"entered_at": r["entered_at"],
|
||||
"exited_at": r["exited_at"],
|
||||
"duration_s": _duration_seconds(r["entered_at"], r["exited_at"]),
|
||||
"status": r["status"],
|
||||
"note": r["note"],
|
||||
"cost_usd": _stage_cost(str(r["to_phase"])),
|
||||
}
|
||||
)
|
||||
elif values:
|
||||
# Best-effort: a single step at the last known phase, no timing.
|
||||
phase = str(values.get("current_phase") or "")
|
||||
node_id = node_for_phase(phase) or phase
|
||||
if node_id:
|
||||
timeline.append(
|
||||
{
|
||||
"from_phase": None,
|
||||
"to_phase": node_id,
|
||||
"entered_at": values.get("created_at"),
|
||||
"exited_at": None,
|
||||
"duration_s": None,
|
||||
"status": values.get("status"),
|
||||
"note": "reconstructed (pre-v3 history)",
|
||||
"cost_usd": _stage_cost(node_id),
|
||||
}
|
||||
)
|
||||
|
||||
return {
|
||||
"ok": True,
|
||||
"thread_id": thread_id,
|
||||
"short_id": status_page._short(thread_id),
|
||||
"task": str(values.get("task") or "(no description)"),
|
||||
"status": str(values.get("status") or "unknown"),
|
||||
"current_phase": str(values.get("current_phase") or "unknown"),
|
||||
"partial": partial,
|
||||
"timeline": timeline,
|
||||
"qa_history": values.get("qa_history") or [],
|
||||
"review_verdicts": values.get("review_verdicts") or [],
|
||||
"plan": values.get("plan"),
|
||||
"cost_by_stage": cost_by_stage,
|
||||
"total_usd": round(sum(cost_by_stage.values()), 6),
|
||||
}
|
||||
|
||||
|
||||
def make_dashboard_app(db_path: str | Path | None = None) -> Any:
|
||||
"""Build the read-only dashboard FastAPI app (does NOT start it).
|
||||
|
||||
``db_path`` defaults to the same ledger the CLI/coordinator use
|
||||
(``AGENT_TEAM_DB`` or the package default). API routes are registered before
|
||||
the static SPA mount so ``/api/*`` always resolves to JSON.
|
||||
"""
|
||||
try:
|
||||
from fastapi import FastAPI, HTTPException
|
||||
from fastapi.responses import JSONResponse
|
||||
except ImportError as exc: # pragma: no cover - depends on optional dep
|
||||
raise RuntimeError(
|
||||
"fastapi is unavailable; install it to serve the dashboard: "
|
||||
"pip install fastapi uvicorn"
|
||||
) from exc
|
||||
|
||||
resolved_db = (
|
||||
Path(db_path) if db_path is not None else status_page._resolve_db_path()
|
||||
)
|
||||
|
||||
app = FastAPI(
|
||||
title="Sea Haven agent-team dashboard",
|
||||
description="Read-only LAN status dashboard for the agent-team pipeline.",
|
||||
version="1.0.0",
|
||||
docs_url=None,
|
||||
redoc_url=None,
|
||||
openapi_url=None,
|
||||
)
|
||||
|
||||
@app.get("/api/state")
|
||||
def api_state() -> JSONResponse:
|
||||
snapshot = status_page.build_snapshot(resolved_db)
|
||||
payload = status_page.snapshot_to_dict(snapshot)
|
||||
payload["nodes"] = _node_states(snapshot) if snapshot.ok else {}
|
||||
return JSONResponse(payload, headers={"Cache-Control": "no-store"})
|
||||
|
||||
@app.get("/api/topology")
|
||||
def api_topology() -> JSONResponse:
|
||||
try:
|
||||
payload = build_topology()
|
||||
except Exception as exc: # pragma: no cover - defensive
|
||||
payload = {"trees": [], "nodes": [], "edges": [], "error": str(exc)}
|
||||
return JSONResponse(payload, headers={"Cache-Control": "no-store"})
|
||||
|
||||
@app.get("/api/task/{thread_id}")
|
||||
def api_task(thread_id: str) -> JSONResponse:
|
||||
if not _THREAD_ID_RE.match(thread_id):
|
||||
raise HTTPException(status_code=400, detail="invalid thread_id")
|
||||
payload = task_detail(resolved_db, thread_id)
|
||||
return JSONResponse(payload, headers={"Cache-Control": "no-store"})
|
||||
|
||||
# Static SPA, mounted LAST so it never shadows the /api routes. Skipped when
|
||||
# no build is present (dev/test) so the API is still reachable.
|
||||
dist = _resolve_dist_dir()
|
||||
if dist.is_dir():
|
||||
from fastapi.staticfiles import StaticFiles
|
||||
|
||||
app.mount("/", StaticFiles(directory=str(dist), html=True), name="spa")
|
||||
|
||||
return app
|
||||
|
||||
|
||||
def serve(
|
||||
*,
|
||||
host: str = _DEFAULT_HOST,
|
||||
port: int = _DEFAULT_PORT,
|
||||
db_path: str | Path | None = None,
|
||||
) -> None: # pragma: no cover - process entry point
|
||||
"""Start the uvicorn server for the read-only dashboard (LAN, :8770).
|
||||
|
||||
This is the systemd ExecStart target replacing ``status_page.serve``.
|
||||
"""
|
||||
try:
|
||||
import uvicorn
|
||||
except ImportError as exc:
|
||||
raise RuntimeError(
|
||||
"uvicorn is unavailable; install it: pip install uvicorn"
|
||||
) from exc
|
||||
uvicorn.run(make_dashboard_app(db_path=db_path), host=host, port=port)
|
||||
|
|
@ -57,14 +57,11 @@ Config (all env, with defaults):
|
|||
|
||||
from __future__ import annotations
|
||||
|
||||
import html
|
||||
import json
|
||||
import os
|
||||
import sqlite3
|
||||
from collections import Counter
|
||||
from dataclasses import asdict, dataclass, field
|
||||
from datetime import datetime, timezone
|
||||
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
|
|
@ -73,18 +70,6 @@ from typing import Any
|
|||
_PACKAGE_DIR = Path(__file__).resolve().parent
|
||||
_DEFAULT_DB = _PACKAGE_DIR.parent / "state" / "agent_team.sqlite"
|
||||
|
||||
_DEFAULT_HOST = "0.0.0.0"
|
||||
_DEFAULT_PORT = 8770
|
||||
|
||||
# <noscript> fallback / hard cap: how long the browser waits before a full
|
||||
# reload if JS is disabled. The live JS poller (below) refreshes far more often
|
||||
# without reloading, so this only matters for the no-JS path.
|
||||
_REFRESH_SECONDS = 10
|
||||
|
||||
# How often the inline JS poller re-fetches /api/state (milliseconds). ~4s keeps
|
||||
# the map feeling live without hammering a read-only SQLite reader.
|
||||
_POLL_INTERVAL_MS = 4000
|
||||
|
||||
# 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"}
|
||||
|
|
@ -587,711 +572,7 @@ def _read_recent_spend(
|
|||
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; }
|
||||
|
||||
/* --- Pipeline map ------------------------------------------------------- */
|
||||
.map-wrap {
|
||||
background: #12151c; border: 1px solid #232a36; border-radius: 10px;
|
||||
padding: 1rem; overflow-x: auto; position: relative;
|
||||
}
|
||||
svg.map { display: block; max-width: 100%; height: auto; }
|
||||
/* Stage node states. Fill = the box; stroke = its border accent. */
|
||||
.node rect.box { fill: #1a1e26; stroke: #2f3645; stroke-width: 1.5;
|
||||
transition: fill .25s, stroke .25s; rx: 9; }
|
||||
.node text { fill: #e6e6e6; font-size: 13px; font-weight: 600;
|
||||
text-anchor: middle; pointer-events: none; }
|
||||
.node text.agent { fill: #8a93a6; font-size: 10.5px; font-weight: 500; }
|
||||
.node text.gated-tag { fill: #6b7385; font-size: 9px; font-weight: 600;
|
||||
letter-spacing: .06em; }
|
||||
.node .count { font-size: 11px; font-weight: 700; fill: #0f1115; }
|
||||
.node circle.badge { display: none; }
|
||||
.node.has-tasks circle.badge { display: inline; }
|
||||
/* state -> accent. idle is the default rect above. */
|
||||
.node.active rect.box { fill: #163a2b; stroke: #2f8f5e; }
|
||||
.node.active circle.badge { fill: #6ee7a8; }
|
||||
.node.awaiting_human rect.box { fill: #3a3216; stroke: #b89a2e; }
|
||||
.node.awaiting_human circle.badge { fill: #f5d76e; }
|
||||
.node.parked rect.box { fill: #3a1f16; stroke: #a85f33; }
|
||||
.node.parked circle.badge { fill: #f3a36e; }
|
||||
.node.gated rect.box { stroke-dasharray: 5 4; opacity: .82; }
|
||||
.node.role rect.box { fill: #181c28; stroke-dasharray: 2 3; }
|
||||
.node:focus { outline: none; }
|
||||
.node:focus rect.box, .node:hover rect.box { stroke-width: 2.5; }
|
||||
.edge { stroke: #3b4456; stroke-width: 1.6; fill: none;
|
||||
marker-end: url(#arrow); }
|
||||
.edge.loop { stroke: #4a5266; stroke-dasharray: 4 4; }
|
||||
.edge-label { fill: #6b7385; font-size: 9.5px; text-anchor: middle; }
|
||||
/* Legend */
|
||||
.legend { display: flex; flex-wrap: wrap; gap: 1rem; margin: .75rem 0 0;
|
||||
font-size: .8rem; color: #9aa3b4; align-items: center; }
|
||||
.legend .key { display: inline-flex; align-items: center; gap: .35rem; }
|
||||
.legend .swatch { width: 14px; height: 14px; border-radius: 4px;
|
||||
border: 1px solid #2f3645; display: inline-block; }
|
||||
.legend .swatch.active { background: #163a2b; border-color: #2f8f5e; }
|
||||
.legend .swatch.awaiting_human { background: #3a3216; border-color: #b89a2e; }
|
||||
.legend .swatch.parked { background: #3a1f16; border-color: #a85f33; }
|
||||
.legend .swatch.idle { background: #1a1e26; }
|
||||
/* Tooltip (one node, follows hover/focus). */
|
||||
#tip {
|
||||
position: fixed; z-index: 10; max-width: 24rem; pointer-events: none;
|
||||
background: #0b0d12; border: 1px solid #36405a; border-radius: 8px;
|
||||
padding: .55rem .7rem; font-size: .8rem; color: #d6dcea;
|
||||
box-shadow: 0 8px 24px rgba(0,0,0,.45); display: none; line-height: 1.35;
|
||||
}
|
||||
#tip.show { display: block; }
|
||||
#tip h4 { margin: 0 0 .35rem; font-size: .82rem; color: #9fb4d6; }
|
||||
#tip .agent-line { color: #8a93a6; font-size: .74rem; margin-bottom: .35rem; }
|
||||
#tip ul { margin: .2rem 0 0; padding-left: 0; list-style: none; }
|
||||
#tip li { margin: .3rem 0; border-top: 1px solid #20283a; padding-top: .3rem; }
|
||||
#tip li:first-child { border-top: 0; padding-top: 0; }
|
||||
#tip .tid { font-family: ui-monospace, SFMono-Regular, Menlo, monospace;
|
||||
color: #7fb3f5; }
|
||||
#tip .empty-tip { color: #6b7385; font-style: italic; }
|
||||
#tip .age { color: #8a93a6; }
|
||||
.live-dot { width: .55rem; height: .55rem; border-radius: 50%;
|
||||
background: #6ee7a8; display: inline-block; margin-right: .3rem;
|
||||
vertical-align: middle; }
|
||||
.live-dot.stale { background: #f3a36e; }
|
||||
@media (prefers-reduced-motion: reduce) {
|
||||
.node rect.box, .edge { transition: none; }
|
||||
}
|
||||
"""
|
||||
|
||||
|
||||
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'<span class="badge {_esc(css)}">{_esc(text)}</span>'
|
||||
|
||||
|
||||
def _render_task_rows(tasks: list[TaskView]) -> str:
|
||||
if not tasks:
|
||||
return '<tr><td colspan="4" class="empty">No tasks.</td></tr>'
|
||||
out: list[str] = []
|
||||
for t in tasks:
|
||||
out.append(
|
||||
"<tr>"
|
||||
f'<td class="mono">{_esc(t.short_id)}</td>'
|
||||
f'<td class="desc">{_esc(t.task)}</td>'
|
||||
f"<td>{_badge(t.current_phase)}</td>"
|
||||
f"<td>{_badge(t.status)}</td>"
|
||||
"</tr>"
|
||||
)
|
||||
return "".join(out)
|
||||
|
||||
|
||||
def _render_waiting(tasks: list[TaskView]) -> str:
|
||||
waiting = [t for t in tasks if t.waiting]
|
||||
if not waiting:
|
||||
return '<p class="empty">No tasks are waiting on the human gate.</p>'
|
||||
rows: list[str] = []
|
||||
for t in waiting:
|
||||
since = _esc(t.waiting_since) if t.waiting_since else "—"
|
||||
rows.append(
|
||||
"<tr>"
|
||||
f'<td class="mono">{_esc(t.short_id)}</td>'
|
||||
f'<td class="desc">{_esc(t.task)}</td>'
|
||||
f"<td>{_badge(t.current_phase)}</td>"
|
||||
f'<td class="mono">{since}</td>'
|
||||
"</tr>"
|
||||
)
|
||||
return (
|
||||
"<table><thead><tr><th>Thread</th><th>Task</th><th>Phase</th>"
|
||||
"<th>Open since (UTC)</th></tr></thead><tbody>"
|
||||
+ "".join(rows)
|
||||
+ "</tbody></table>"
|
||||
)
|
||||
|
||||
|
||||
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(
|
||||
"<tr>"
|
||||
f'<td class="mono">{_esc(entry.get("recorded_at"))}</td>'
|
||||
f"<td>{_esc(entry.get('stage') or '—')}</td>"
|
||||
f'<td class="mono">{_esc(entry.get("model"))}</td>'
|
||||
f'<td class="mono">${float(entry.get("usd_cost") or 0.0):,.4f}</td>'
|
||||
"</tr>"
|
||||
)
|
||||
body = (
|
||||
"".join(rows)
|
||||
if rows
|
||||
else '<tr><td colspan="4" class="empty">No spend recorded.</td></tr>'
|
||||
)
|
||||
return (
|
||||
f"<h2>Budget — total recorded spend {total}</h2>"
|
||||
"<table><thead><tr><th>Recorded (UTC)</th><th>Stage</th>"
|
||||
"<th>Model</th><th>USD</th></tr></thead><tbody>" + body + "</tbody></table>"
|
||||
)
|
||||
|
||||
|
||||
# --- Inline SVG pipeline map (hand-rolled, no CDN, no D3). -------------------
|
||||
#
|
||||
# Layout: a serpentine of two rows of five stage nodes, flowing row 1 L->R then
|
||||
# row 2 L->R (with a connector down the right side), so the whole DAG reads in
|
||||
# one screen. Loop-backs (CLARIFY <-> human gate, PLAN <-> REVIEW) are drawn as
|
||||
# dashed curved edges so the operator can see the gating cycles.
|
||||
|
||||
_NODE_W = 132
|
||||
_NODE_H = 66
|
||||
_COL_GAP = 56
|
||||
_ROW_GAP = 86
|
||||
_MARGIN = 24
|
||||
_COLS = 5
|
||||
|
||||
|
||||
def _node_xy(index: int) -> tuple[int, int]:
|
||||
"""Top-left (x, y) of the stage box at ``index`` in the serpentine grid."""
|
||||
row, col = divmod(index, _COLS)
|
||||
x = _MARGIN + col * (_NODE_W + _COL_GAP)
|
||||
y = _MARGIN + row * (_NODE_H + _ROW_GAP)
|
||||
return x, y
|
||||
|
||||
|
||||
def _render_node(index: int, stage: Stage, state: str, count: int) -> str:
|
||||
"""One stage box: state class, label, agent role, gated tag, count badge."""
|
||||
x, y = _node_xy(index)
|
||||
classes = ["node", _esc(state)]
|
||||
if stage.gated:
|
||||
classes.append("gated")
|
||||
if stage.role_node:
|
||||
classes.append("role")
|
||||
if count > 0:
|
||||
classes.append("has-tasks")
|
||||
cls = " ".join(classes)
|
||||
cx, cy = x + _NODE_W, y # badge anchor: top-right corner
|
||||
label_x = x + _NODE_W // 2
|
||||
gated_tag = (
|
||||
f'<text class="gated-tag" x="{label_x}" y="{y + _NODE_H - 8}">GATED</text>'
|
||||
if stage.gated
|
||||
else ""
|
||||
)
|
||||
badge = (
|
||||
f'<circle class="badge" cx="{cx}" cy="{cy}" r="11"></circle>'
|
||||
f'<text class="count" x="{cx}" y="{cy + 4}">{count}</text>'
|
||||
)
|
||||
return (
|
||||
f'<g class="{cls}" tabindex="0" role="button" '
|
||||
f'data-stage="{_esc(stage.key)}" '
|
||||
f'aria-label="{_esc(stage.label)} stage, {_esc(state)}, '
|
||||
f'{count} task(s)">'
|
||||
f'<rect class="box" x="{x}" y="{y}" width="{_NODE_W}" '
|
||||
f'height="{_NODE_H}" rx="9"></rect>'
|
||||
f'<text x="{label_x}" y="{y + 27}">{_esc(stage.label)}</text>'
|
||||
f'<text class="agent" x="{label_x}" y="{y + 45}">'
|
||||
f"{_esc(stage.agent)}</text>"
|
||||
f"{gated_tag}{badge}</g>"
|
||||
)
|
||||
|
||||
|
||||
def _straight_edge(a: int, b: int) -> str:
|
||||
"""Forward edge from node ``a`` to node ``b`` (adjacent in the spine)."""
|
||||
ax, ay = _node_xy(a)
|
||||
bx, by = _node_xy(b)
|
||||
a_row = a // _COLS
|
||||
b_row = b // _COLS
|
||||
if a_row == b_row:
|
||||
# Same row: right edge of a -> left edge of b.
|
||||
x1, y1 = ax + _NODE_W, ay + _NODE_H // 2
|
||||
x2, y2 = bx, by + _NODE_H // 2
|
||||
mx = (x1 + x2) / 2
|
||||
return (
|
||||
f'<path class="edge" d="M{x1},{y1} C{mx},{y1} {mx},{y2} {x2},{y2}"></path>'
|
||||
)
|
||||
# Row wrap: drop from bottom of a, curve down to top of b.
|
||||
x1, y1 = ax + _NODE_W // 2, ay + _NODE_H
|
||||
x2, y2 = bx + _NODE_W // 2, by
|
||||
my = (y1 + y2) / 2
|
||||
return f'<path class="edge" d="M{x1},{y1} C{x1},{my} {x2},{my} {x2},{y2}"></path>'
|
||||
|
||||
|
||||
def _loop_edge(a: int, b: int, label: str) -> str:
|
||||
"""Dashed loop-back edge a<->b drawn as an arc above the two same-row nodes."""
|
||||
ax, ay = _node_xy(a)
|
||||
bx, _ = _node_xy(b)
|
||||
# Arc over the top of both boxes (a is to the right of b on the spine).
|
||||
x1, y1 = ax + _NODE_W // 2, ay
|
||||
x2, y2 = bx + _NODE_W // 2, ay
|
||||
lift = 30
|
||||
mid_x = (x1 + x2) / 2
|
||||
path = (
|
||||
f'<path class="edge loop" d="M{x1},{y1} '
|
||||
f'C{x1},{y1 - lift} {x2},{y2 - lift} {x2},{y2}"></path>'
|
||||
)
|
||||
lbl = (
|
||||
f'<text class="edge-label" x="{mid_x}" y="{y1 - lift + 2}">{_esc(label)}</text>'
|
||||
)
|
||||
return path + lbl
|
||||
|
||||
|
||||
def _render_map_svg(snapshot: Snapshot) -> str:
|
||||
"""Hand-rolled inline SVG of the agent-team DAG, with live per-stage state.
|
||||
|
||||
Deterministic geometry (no measurement), all inline — no external assets.
|
||||
The JS poller re-paints node classes/counts in place; this server render is
|
||||
the first paint and the no-JS view.
|
||||
"""
|
||||
rows = (len(STAGES) + _COLS - 1) // _COLS
|
||||
width = _MARGIN * 2 + _COLS * _NODE_W + (_COLS - 1) * _COL_GAP
|
||||
height = _MARGIN * 2 + rows * _NODE_H + (rows - 1) * _ROW_GAP
|
||||
|
||||
parts: list[str] = []
|
||||
# Forward spine: 0->1->...->n following serpentine order.
|
||||
for i in range(len(STAGES) - 1):
|
||||
parts.append(_straight_edge(i, i + 1))
|
||||
# Loop-backs by stage key (robust to index shifts).
|
||||
idx = {s.key: i for i, s in enumerate(STAGES)}
|
||||
if "clarify" in idx and "gate" in idx:
|
||||
parts.append(_loop_edge(idx["gate"], idx["clarify"], "answer ↺"))
|
||||
if "plan" in idx and "review" in idx:
|
||||
parts.append(_loop_edge(idx["review"], idx["plan"], "revise ↺"))
|
||||
|
||||
for i, stage in enumerate(STAGES):
|
||||
members = snapshot.tasks_for_stage(stage.key)
|
||||
state = snapshot.stage_state(stage)
|
||||
parts.append(_render_node(i, stage, state, len(members)))
|
||||
|
||||
arrow = (
|
||||
'<defs><marker id="arrow" viewBox="0 0 10 10" refX="9" refY="5" '
|
||||
'markerWidth="7" markerHeight="7" orient="auto-start-reverse">'
|
||||
'<path d="M0,0 L10,5 L0,10 z" fill="#3b4456"></path></marker></defs>'
|
||||
)
|
||||
return (
|
||||
f'<svg class="map" viewBox="0 0 {width} {height}" '
|
||||
f'width="{width}" height="{height}" role="img" '
|
||||
'aria-label="agent-team pipeline map">'
|
||||
f"{arrow}{''.join(parts)}</svg>"
|
||||
)
|
||||
|
||||
|
||||
def _render_legend() -> str:
|
||||
return (
|
||||
'<div class="legend">'
|
||||
'<span class="key"><span class="swatch idle"></span>idle</span>'
|
||||
'<span class="key"><span class="swatch active"></span>active</span>'
|
||||
'<span class="key"><span class="swatch awaiting_human"></span>'
|
||||
"awaiting human</span>"
|
||||
'<span class="key"><span class="swatch parked"></span>parked</span>'
|
||||
'<span class="key" style="color:#6b7385">'
|
||||
"dashed border = gated / not-yet-live stage</span>"
|
||||
"</div>"
|
||||
)
|
||||
|
||||
|
||||
def _json_for_script(payload: dict[str, Any]) -> str:
|
||||
"""JSON safe to inline inside a ``<script>`` element.
|
||||
|
||||
``json.dumps`` handles the JSON-string escaping. We then escape ``<`` and
|
||||
``>`` to their ``\\uXXXX`` JSON escapes, which is the standard hardening for
|
||||
JSON embedded in HTML: it makes it impossible for any value (e.g. a task
|
||||
description containing ``</script>`` or ``<script>``) to introduce a tag,
|
||||
comment, or CDATA breakout, while remaining byte-for-byte parseable by
|
||||
``JSON.parse``. ``ensure_ascii=True`` additionally escapes the unicode line
|
||||
separators that would otherwise terminate a JS string. Descriptions still
|
||||
only ever reach the DOM via ``textContent``, never ``innerHTML``.
|
||||
"""
|
||||
raw = json.dumps(payload, ensure_ascii=True)
|
||||
return raw.replace("<", "\\u003c").replace(">", "\\u003e")
|
||||
|
||||
|
||||
def _poller_script(payload: dict[str, Any]) -> str:
|
||||
"""Inline vanilla-JS poller: fetch /api/state, repaint the map in place.
|
||||
|
||||
No libraries, no CDN. Updates node state classes + count badges, the count
|
||||
cards, the "last updated" clock, and the tooltip dataset — without reloading
|
||||
the page, so hover/scroll/focus survive. Descriptions reach the DOM only via
|
||||
``textContent`` (never innerHTML), so they cannot inject markup.
|
||||
"""
|
||||
seed = _json_for_script(payload)
|
||||
return (
|
||||
"<script>(function(){"
|
||||
"var POLL=" + str(_POLL_INTERVAL_MS) + ";"
|
||||
"var state=" + seed + ";"
|
||||
"var tip=document.getElementById('tip');"
|
||||
"var clock=document.getElementById('clock');"
|
||||
"var dot=document.getElementById('livedot');"
|
||||
# ---- helpers ----
|
||||
"function age(iso){if(!iso)return '';"
|
||||
"var t=Date.parse(iso);if(isNaN(t))return '';"
|
||||
"var s=Math.max(0,(Date.now()-t)/1000);"
|
||||
"if(s<90)return Math.round(s)+'s ago';"
|
||||
"if(s<5400)return Math.round(s/60)+'m ago';"
|
||||
"if(s<129600)return Math.round(s/3600)+'h ago';"
|
||||
"return Math.round(s/86400)+'d ago';}"
|
||||
"function byKey(k){var s=state.stages||[];"
|
||||
"for(var i=0;i<s.length;i++){if(s[i].key===k)return s[i];}return null;}"
|
||||
"function setNum(id,v){var e=document.getElementById(id);"
|
||||
"if(e)e.textContent=String(v==null?0:v);}"
|
||||
# ---- repaint the SVG nodes + cards from `state` ----
|
||||
"function paint(){"
|
||||
"var nodes=document.querySelectorAll('g.node[data-stage]');"
|
||||
"for(var i=0;i<nodes.length;i++){var g=nodes[i];"
|
||||
"var st=byKey(g.getAttribute('data-stage'));if(!st)continue;"
|
||||
"['idle','active','awaiting_human','parked','has-tasks']"
|
||||
".forEach(function(c){g.classList.remove(c);});"
|
||||
"g.classList.add(st.state||'idle');"
|
||||
"if((st.count||0)>0)g.classList.add('has-tasks');"
|
||||
"var cnt=g.querySelector('text.count');"
|
||||
"if(cnt)cnt.textContent=String(st.count||0);"
|
||||
"g.setAttribute('aria-label',(st.label||g.getAttribute('data-stage'))"
|
||||
"+' stage, '+(st.state||'idle')+', '+(st.count||0)+' task(s)');}"
|
||||
"var sm=state.summary||{};"
|
||||
"setNum('card-active',sm.active);setNum('card-parked',sm.parked);"
|
||||
"setNum('card-waiting',sm.waiting);setNum('card-open',sm.open_questions);"
|
||||
"setNum('card-total',sm.total);}"
|
||||
# ---- tooltip ----
|
||||
"function fillTip(st){"
|
||||
"var h=document.createElement('h4');h.textContent=(st.label||'')+' — '"
|
||||
"+(st.state||'idle').replace('_',' ');"
|
||||
"var a=document.createElement('div');a.className='agent-line';"
|
||||
"a.textContent='agent: '+(st.agent||'—');"
|
||||
"tip.innerHTML='';tip.appendChild(h);tip.appendChild(a);"
|
||||
"var tasks=st.tasks||[];"
|
||||
"if(!tasks.length){var e=document.createElement('div');"
|
||||
"e.className='empty-tip';"
|
||||
"e.textContent=(st.gated?'gated / not yet live — ':'')+'no tasks here';"
|
||||
"tip.appendChild(e);return;}"
|
||||
"var ul=document.createElement('ul');"
|
||||
"for(var i=0;i<tasks.length;i++){var t=tasks[i];"
|
||||
"var li=document.createElement('li');"
|
||||
"var idline=document.createElement('div');"
|
||||
"var sp=document.createElement('span');sp.className='tid';"
|
||||
"sp.textContent=t.short_id||t.thread_id;"
|
||||
"idline.appendChild(sp);"
|
||||
"idline.appendChild(document.createTextNode(' · '+(t.status||'')));"
|
||||
"li.appendChild(idline);"
|
||||
"var d=document.createElement('div');d.textContent=t.task||'(no description)';"
|
||||
"li.appendChild(d);"
|
||||
"var when=t.waiting?t.waiting_since:null;"
|
||||
"if(when){var ag=document.createElement('div');ag.className='age';"
|
||||
"ag.textContent='waiting since '+when+(age(when)?' ('+age(when)+')':'');"
|
||||
"li.appendChild(ag);}"
|
||||
"ul.appendChild(li);}"
|
||||
"tip.appendChild(ul);}"
|
||||
"function showTip(key,ev){var st=byKey(key);if(!st)return;"
|
||||
"fillTip(st);tip.classList.add('show');moveTip(ev);}"
|
||||
"function moveTip(ev){if(!ev)return;var pad=14;"
|
||||
"var x=ev.clientX+pad,y=ev.clientY+pad;"
|
||||
"var r=tip.getBoundingClientRect();"
|
||||
"if(x+r.width>window.innerWidth)x=ev.clientX-r.width-pad;"
|
||||
"if(y+r.height>window.innerHeight)y=ev.clientY-r.height-pad;"
|
||||
"tip.style.left=x+'px';tip.style.top=y+'px';}"
|
||||
"function hideTip(){tip.classList.remove('show');}"
|
||||
# ---- wire node hover/focus ----
|
||||
"function wire(){var nodes=document.querySelectorAll('g.node[data-stage]');"
|
||||
"for(var i=0;i<nodes.length;i++){(function(g){"
|
||||
"var k=g.getAttribute('data-stage');"
|
||||
"g.addEventListener('mouseenter',function(e){showTip(k,e);});"
|
||||
"g.addEventListener('mousemove',function(e){moveTip(e);});"
|
||||
"g.addEventListener('mouseleave',hideTip);"
|
||||
"g.addEventListener('focus',function(){var st=byKey(k);if(st){"
|
||||
"fillTip(st);var b=g.getBoundingClientRect();tip.classList.add('show');"
|
||||
"tip.style.left=(b.left)+'px';tip.style.top=(b.bottom+8)+'px';}});"
|
||||
"g.addEventListener('blur',hideTip);"
|
||||
"})(nodes[i]);}}"
|
||||
# ---- clock ----
|
||||
"function setClock(){if(clock)clock.textContent=state.generated_at||'';"
|
||||
"if(dot){dot.classList.remove('stale');}}"
|
||||
# ---- poll ----
|
||||
"function apply(s){state=s;paint();setClock();}"
|
||||
"function poll(){fetch('/api/state',{cache:'no-store'})"
|
||||
".then(function(r){return r.json();})"
|
||||
".then(function(s){apply(s);})"
|
||||
".catch(function(){if(dot)dot.classList.add('stale');});}"
|
||||
# ---- init ----
|
||||
"paint();wire();setClock();"
|
||||
"setInterval(poll,POLL);"
|
||||
"})();</script>"
|
||||
)
|
||||
|
||||
|
||||
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 + inline SVG map + inline JS poller only, plus a
|
||||
``<noscript>`` meta-refresh fallback, so the page works fully offline (no
|
||||
external CDN) and updates live without a full reload.
|
||||
"""
|
||||
payload = snapshot_to_dict(snapshot)
|
||||
head = (
|
||||
'<!doctype html><html lang="en"><head><meta charset="utf-8">'
|
||||
'<meta name="viewport" content="width=device-width, initial-scale=1">'
|
||||
# No-JS fallback: full reload. The inline poller refreshes live without
|
||||
# reloading, so this only fires when JS is disabled.
|
||||
"<noscript>"
|
||||
f'<meta http-equiv="refresh" content="{_REFRESH_SECONDS}">'
|
||||
"</noscript>"
|
||||
"<title>agent-team status</title>"
|
||||
f"<style>{_CSS}</style></head><body>"
|
||||
)
|
||||
header = (
|
||||
"<h1>Sea Haven agent-team — coordinator status</h1>"
|
||||
'<div class="meta"><span class="live-dot" id="livedot"></span>'
|
||||
'last updated <span class="mono" id="clock">'
|
||||
f"{_esc(snapshot.generated_at)}</span> · live "
|
||||
f"(polls /api/state every {_POLL_INTERVAL_MS // 1000}s) · "
|
||||
f'<span class="mono">{_esc(snapshot.db_path)}</span> · '
|
||||
"read-only, LAN/VPN-only</div>"
|
||||
)
|
||||
tooltip = '<div id="tip" role="tooltip"></div>'
|
||||
|
||||
# The map + tooltip + poller render on every path (including no-data) so the
|
||||
# page goes live the moment the coordinator produces state, without a reload.
|
||||
map_section = (
|
||||
"<h2>Pipeline map</h2>"
|
||||
'<div class="map-wrap">'
|
||||
+ _render_map_svg(snapshot)
|
||||
+ "</div>"
|
||||
+ _render_legend()
|
||||
)
|
||||
|
||||
if not snapshot.ok:
|
||||
warn = (
|
||||
'<div class="warn"><strong>No data.</strong><br>'
|
||||
f"{_esc(snapshot.error or 'ledger unavailable')}<br>"
|
||||
"The dashboard is read-only and waits for the coordinator to "
|
||||
"produce state. The map will populate live once it does.</div>"
|
||||
)
|
||||
return (
|
||||
head
|
||||
+ header
|
||||
+ tooltip
|
||||
+ map_section
|
||||
+ warn
|
||||
+ _footer()
|
||||
+ _poller_script(payload)
|
||||
+ "</body></html>"
|
||||
)
|
||||
|
||||
qc = snapshot.question_counts
|
||||
cards = (
|
||||
'<div class="cards">'
|
||||
f'<div class="card"><div class="num" id="card-active">'
|
||||
f"{snapshot.active_count}</div>"
|
||||
'<div class="lbl">Active</div></div>'
|
||||
f'<div class="card"><div class="num" id="card-parked">'
|
||||
f"{snapshot.parked_count}</div>"
|
||||
'<div class="lbl">Parked</div></div>'
|
||||
f'<div class="card"><div class="num" id="card-waiting">'
|
||||
f"{len(snapshot.waiting_tasks)}</div>"
|
||||
'<div class="lbl">Waiting (gate)</div></div>'
|
||||
f'<div class="card"><div class="num" id="card-open">'
|
||||
f"{qc.get('open', 0)}</div>"
|
||||
'<div class="lbl">Open questions</div></div>'
|
||||
f'<div class="card"><div class="num" id="card-total">'
|
||||
f"{len(snapshot.tasks)}</div>"
|
||||
'<div class="lbl">Total tasks</div></div>'
|
||||
"</div>"
|
||||
)
|
||||
|
||||
waiting_section = "<h2>Waiting on the human gate</h2>" + _render_waiting(
|
||||
snapshot.tasks
|
||||
)
|
||||
|
||||
tasks_section = (
|
||||
"<h2>All tasks</h2>"
|
||||
"<table><thead><tr><th>Thread</th><th>Task</th><th>Phase</th>"
|
||||
"<th>Status</th></tr></thead><tbody>"
|
||||
+ _render_task_rows(snapshot.tasks)
|
||||
+ "</tbody></table>"
|
||||
)
|
||||
|
||||
budget_section = _render_budget(snapshot)
|
||||
|
||||
return (
|
||||
head
|
||||
+ header
|
||||
+ tooltip
|
||||
+ cards
|
||||
+ map_section
|
||||
+ waiting_section
|
||||
+ tasks_section
|
||||
+ budget_section
|
||||
+ _footer()
|
||||
+ _poller_script(payload)
|
||||
+ "</body></html>"
|
||||
)
|
||||
|
||||
|
||||
def _footer() -> str:
|
||||
return (
|
||||
"<footer>READ-ONLY status dashboard · no mutating endpoints "
|
||||
"· unauthenticated — do not expose to the public "
|
||||
"internet. Task descriptions may be sensitive.</footer>"
|
||||
)
|
||||
|
||||
|
||||
# --- 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"
|
||||
|
||||
# Strip any query string so "/api/state?ts=123" (a cache-buster the JS
|
||||
# may append) still routes. Read-only: only GET, only these two routes.
|
||||
def _route(self) -> str:
|
||||
return self.path.split("?", 1)[0]
|
||||
|
||||
def do_GET(self) -> None: # noqa: N802 (BaseHTTPRequestHandler API)
|
||||
route = self._route()
|
||||
# Favicon gets a 204 so browsers stop logging 404s.
|
||||
if route == "/favicon.ico":
|
||||
self.send_response(204)
|
||||
self.end_headers()
|
||||
return
|
||||
if route == "/api/state":
|
||||
self._serve_json()
|
||||
return
|
||||
# Every other path serves the dashboard HTML.
|
||||
self._serve_html()
|
||||
|
||||
def _serve_json(self) -> None:
|
||||
try:
|
||||
payload = snapshot_to_dict(build_snapshot(db_path))
|
||||
except Exception as exc: # last-ditch: never 500 the sidecar
|
||||
payload = {
|
||||
"ok": False,
|
||||
"error": f"snapshot error: {exc}",
|
||||
"generated_at": _utc_now_iso(),
|
||||
"db_path": str(db_path),
|
||||
"stages": [],
|
||||
"tasks": [],
|
||||
"waiting": [],
|
||||
"question_counts": {},
|
||||
"summary": {},
|
||||
"budget": {"available": False},
|
||||
}
|
||||
body = json.dumps(payload).encode("utf-8")
|
||||
self.send_response(200)
|
||||
self.send_header("Content-Type", "application/json; charset=utf-8")
|
||||
self.send_header("Content-Length", str(len(body)))
|
||||
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 _serve_html(self) -> None:
|
||||
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()
|
||||
# NOTE: HTML/SVG rendering and the stdlib http.server were retired in the WebUI
|
||||
# makeover. The dashboard is now a React/Vite SPA served by agent_team.dashboard
|
||||
# (FastAPI). This module is kept as the read-only DATA LAYER: build_snapshot()
|
||||
# and snapshot_to_dict() feed the dashboard's /api/state endpoint.
|
||||
|
|
|
|||
260
agent-team/agent_team/topology.py
Normal file
260
agent-team/agent_team/topology.py
Normal file
|
|
@ -0,0 +1,260 @@
|
|||
"""Pipeline topology for the status dashboard, derived from the real graph.
|
||||
|
||||
The map the WebUI draws is **introspected from the compiled LangGraph** rather
|
||||
than hand-laid: :func:`build_topology` assembles the maximal graph wiring (all
|
||||
optional P2/P3 nodes injected as lightweight stubs — we want the *shape*, not
|
||||
the behaviour), calls ``compiled.get_graph()`` for its nodes + edges, and merges
|
||||
that with :data:`NODE_META`, a display-only sidecar (label / owning agent /
|
||||
process tree / kind / gated).
|
||||
|
||||
Why derive instead of hardcode: adding a new agent node in
|
||||
:mod:`agent_team.graph` makes it appear on the map automatically — the
|
||||
"easy to add nodes" goal. ``NODE_META`` supplies only presentation; a graph node
|
||||
missing from it still renders (raw id, ``unassigned`` tree) so a new agent is
|
||||
never silently dropped, and :func:`missing_meta` lets a test fail loudly until
|
||||
its metadata is filled in.
|
||||
|
||||
``tree`` groups nodes into processes branching off the ``intake`` coordinator
|
||||
(today: the ``core`` root + the ``sdlc`` pipeline); new processes are new
|
||||
``tree`` values across their nodes' meta entries.
|
||||
|
||||
This module is import-light and has no live model / DB dependency: the stub
|
||||
nodes/routes are never executed (``get_graph()`` reads the wiring statically), so
|
||||
building the topology is a pure, cheap, deterministic operation.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections import deque
|
||||
from typing import Any
|
||||
|
||||
from agent_team import graph as graph_mod
|
||||
from agent_team.task_model import Phase
|
||||
|
||||
__all__ = [
|
||||
"NODE_META",
|
||||
"PHASE_TO_NODE",
|
||||
"TREES",
|
||||
"build_topology",
|
||||
"missing_meta",
|
||||
"node_for_phase",
|
||||
]
|
||||
|
||||
# LangGraph's synthetic terminal vertices — excluded from the drawn map.
|
||||
_START = "__start__"
|
||||
_END = "__end__"
|
||||
|
||||
# Display sidecar: graph node id -> presentation. ``kind`` is "phase" for work
|
||||
# stages and "gate" for the clarifier (which holds the human interrupt). ``tree``
|
||||
# groups nodes into processes branching off intake. ``gated`` marks nodes that
|
||||
# are drawn but inert in the default production deploy (the P3 build->verify->
|
||||
# dispatch path), so the UI can dim them. Keep ids in sync with agent_team.graph.
|
||||
NODE_META: dict[str, dict[str, Any]] = {
|
||||
graph_mod.INTAKE: {
|
||||
"label": "Intake",
|
||||
"agent": "coordinator",
|
||||
"tree": "core",
|
||||
"kind": "phase",
|
||||
"gated": False,
|
||||
},
|
||||
graph_mod.CLARIFY: {
|
||||
"label": "Clarify",
|
||||
"agent": "Claude (sub)",
|
||||
"tree": "sdlc",
|
||||
"kind": "gate", # the human gate (LangGraph interrupt) lives here
|
||||
"gated": False,
|
||||
},
|
||||
graph_mod.PLAN: {
|
||||
"label": "Plan",
|
||||
"agent": "Claude (sub)",
|
||||
"tree": "sdlc",
|
||||
"kind": "phase",
|
||||
"gated": False,
|
||||
},
|
||||
graph_mod.REVIEW: {
|
||||
"label": "Review",
|
||||
"agent": "GPT-4.1 (cross)",
|
||||
"tree": "sdlc",
|
||||
"kind": "phase",
|
||||
"gated": False,
|
||||
},
|
||||
graph_mod.BUILD_NODE: {
|
||||
"label": "Build",
|
||||
"agent": "DeepSeek (fast)",
|
||||
"tree": "sdlc",
|
||||
"kind": "phase",
|
||||
"gated": True, # P3, opt-in/inert in the default deploy
|
||||
},
|
||||
graph_mod.VERIFY_NODE: {
|
||||
"label": "Verify",
|
||||
"agent": "Claude (sub)",
|
||||
"tree": "sdlc",
|
||||
"kind": "phase",
|
||||
"gated": True,
|
||||
},
|
||||
graph_mod.DISPATCH_NODE: {
|
||||
"label": "Dispatch",
|
||||
"agent": "GitHub PR",
|
||||
"tree": "sdlc",
|
||||
"kind": "phase",
|
||||
"gated": True,
|
||||
},
|
||||
}
|
||||
|
||||
# Process trees, in display order. ``root`` flags the tree that owns intake (the
|
||||
# branch point); other trees hang off it. New processes append here.
|
||||
TREES: tuple[dict[str, Any], ...] = (
|
||||
{"id": "core", "label": "Core", "root": True},
|
||||
{"id": "sdlc", "label": "SDLC Pipeline", "root": False},
|
||||
)
|
||||
|
||||
# Default tree/meta for a graph node with no NODE_META entry (a newly added
|
||||
# agent whose metadata has not been filled in yet) — rendered, never dropped.
|
||||
_UNASSIGNED_META: dict[str, Any] = {
|
||||
"label": "", # filled with the raw id at build time
|
||||
"agent": "",
|
||||
"tree": "unassigned",
|
||||
"kind": "phase",
|
||||
"gated": False,
|
||||
}
|
||||
|
||||
# Phase value (TaskStatus current_phase) -> graph node id, so /api/state can
|
||||
# bucket a live task onto its node. BUILD/VERIFY phases map to the P3 vertex ids
|
||||
# (build_node/verify_node). DONE/PARKED are terminal/exception states with no
|
||||
# vertex — a task in them is shown in the list, not on a node.
|
||||
PHASE_TO_NODE: dict[str, str] = {
|
||||
Phase.INTAKE.value: graph_mod.INTAKE,
|
||||
Phase.CLARIFY.value: graph_mod.CLARIFY,
|
||||
Phase.PLAN.value: graph_mod.PLAN,
|
||||
Phase.REVIEW.value: graph_mod.REVIEW,
|
||||
Phase.BUILD.value: graph_mod.BUILD_NODE,
|
||||
Phase.VERIFY.value: graph_mod.VERIFY_NODE,
|
||||
}
|
||||
|
||||
|
||||
def node_for_phase(phase: str | None) -> str | None:
|
||||
"""Return the graph node id a live ``current_phase`` belongs to, or ``None``."""
|
||||
if not phase:
|
||||
return None
|
||||
return PHASE_TO_NODE.get(phase)
|
||||
|
||||
|
||||
def _stub_node(state: Any) -> dict[str, Any]:
|
||||
"""A no-op node used only to assemble the maximal graph shape (never run)."""
|
||||
return {}
|
||||
|
||||
|
||||
def _stub_route(state: Any) -> str:
|
||||
"""A no-op router; never executed — get_graph() reads the edge map statically."""
|
||||
return graph_mod.APPROVED_ROUTE
|
||||
|
||||
|
||||
def _maximal_compiled() -> Any:
|
||||
"""Compile the full P3+ wiring with stub nodes, for shape introspection."""
|
||||
return graph_mod.build_graph(
|
||||
None,
|
||||
live_clarify_node=_stub_node,
|
||||
live_plan_node=_stub_node,
|
||||
review_node=_stub_node,
|
||||
route_review=_stub_route,
|
||||
build_verify=(_stub_node, _stub_node, _stub_route),
|
||||
dispatch_node=_stub_node,
|
||||
)
|
||||
|
||||
|
||||
def _bfs_order(node_ids: list[str], edges: list[tuple[str, str]]) -> dict[str, int]:
|
||||
"""Assign each node a forward rank via BFS from ``__start__``.
|
||||
|
||||
Used to classify edges: an edge whose target ranks at or before its source
|
||||
goes backward (a retry loop). Nodes unreachable in the BFS get a large rank
|
||||
so they never spuriously read as loop targets.
|
||||
"""
|
||||
adj: dict[str, list[str]] = {}
|
||||
for src, dst in edges:
|
||||
adj.setdefault(src, []).append(dst)
|
||||
order: dict[str, int] = {}
|
||||
queue: deque[str] = deque([_START])
|
||||
order[_START] = 0
|
||||
while queue:
|
||||
node = queue.popleft()
|
||||
for nxt in adj.get(node, []):
|
||||
if nxt not in order:
|
||||
order[nxt] = order[node] + 1
|
||||
queue.append(nxt)
|
||||
big = len(node_ids) + len(edges) + 1
|
||||
for nid in node_ids:
|
||||
order.setdefault(nid, big)
|
||||
return order
|
||||
|
||||
|
||||
def _edge_kind(src: str, dst: str, conditional: bool, order: dict[str, int]) -> str:
|
||||
"""Classify an edge as spine | branch | loopback for styling."""
|
||||
if conditional and order.get(dst, 0) <= order.get(src, 0):
|
||||
return "loopback"
|
||||
if conditional:
|
||||
return "branch"
|
||||
return "spine"
|
||||
|
||||
|
||||
def missing_meta() -> list[str]:
|
||||
"""Return graph node ids (excluding start/end) that lack a NODE_META entry.
|
||||
|
||||
A test asserts this is empty so a newly added agent fails loudly until its
|
||||
display metadata is filled in (the node still renders meanwhile).
|
||||
"""
|
||||
compiled = _maximal_compiled()
|
||||
drawable = compiled.get_graph()
|
||||
ids = [n for n in drawable.nodes if n not in (_START, _END)]
|
||||
return [nid for nid in ids if nid not in NODE_META]
|
||||
|
||||
|
||||
def build_topology() -> dict[str, Any]:
|
||||
"""Return the dashboard topology: ``{trees, nodes, edges}``.
|
||||
|
||||
Nodes/edges are derived from the compiled LangGraph (so new graph nodes
|
||||
appear automatically) and enriched with :data:`NODE_META`. Edges are
|
||||
de-duplicated and classified spine/branch/loopback. The synthetic
|
||||
``__start__``/``__end__`` vertices are dropped; an edge touching them is
|
||||
dropped too (the map shows agent nodes, not the framework terminals).
|
||||
"""
|
||||
compiled = _maximal_compiled()
|
||||
drawable = compiled.get_graph()
|
||||
|
||||
raw_edges = [(e.source, e.target, bool(e.conditional)) for e in drawable.edges]
|
||||
order = _bfs_order(list(drawable.nodes), [(s, t) for s, t, _ in raw_edges])
|
||||
|
||||
node_ids = [n for n in drawable.nodes if n not in (_START, _END)]
|
||||
nodes: list[dict[str, Any]] = []
|
||||
for nid in node_ids:
|
||||
meta = NODE_META.get(nid)
|
||||
if meta is None:
|
||||
meta = {**_UNASSIGNED_META, "label": nid}
|
||||
nodes.append(
|
||||
{
|
||||
"id": nid,
|
||||
"label": meta["label"] or nid,
|
||||
"agent": meta["agent"],
|
||||
"tree": meta["tree"],
|
||||
"kind": meta["kind"],
|
||||
"gated": bool(meta["gated"]),
|
||||
}
|
||||
)
|
||||
|
||||
seen: set[tuple[str, str]] = set()
|
||||
edges: list[dict[str, Any]] = []
|
||||
for src, dst, conditional in raw_edges:
|
||||
if src in (_START, _END) or dst in (_START, _END):
|
||||
continue
|
||||
key = (src, dst)
|
||||
if key in seen:
|
||||
continue
|
||||
seen.add(key)
|
||||
edges.append(
|
||||
{"from": src, "to": dst, "kind": _edge_kind(src, dst, conditional, order)}
|
||||
)
|
||||
|
||||
return {
|
||||
"trees": [dict(t) for t in TREES],
|
||||
"nodes": nodes,
|
||||
"edges": edges,
|
||||
}
|
||||
Reference in a new issue