diff --git a/agent-team/agent_team/dashboard.py b/agent-team/agent_team/dashboard.py new file mode 100644 index 0000000..feaa38c --- /dev/null +++ b/agent-team/agent_team/dashboard.py @@ -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) diff --git a/agent-team/agent_team/status_page.py b/agent-team/agent_team/status_page.py index 4842949..1d57d78 100644 --- a/agent-team/agent_team/status_page.py +++ b/agent-team/agent_team/status_page.py @@ -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 - -#