feat(agent-team): WebUI makeover — branching pipeline map + click-through task history #55

Merged
amoussa1229 merged 7 commits from feature/agent-team-webui-makeover into main 2026-06-23 23:50:17 +00:00
35 changed files with 6889 additions and 929 deletions

View file

@ -43,10 +43,18 @@ agent-team/
# Gemini via the orchestrator's models.py); bind_multi_invoker()
api.py # WS1 FastAPI HTTP API (bearer auth, 127.0.0.1:8765) — SEPARATE
# opt-in process (api.serve()), NOT started by the coordinator
dashboard.py # read-only LAN status dashboard (FastAPI, 0.0.0.0:8770):
# /api/state /api/topology /api/task/{id} + serves web/dist SPA
topology.py # pipeline map derived from the compiled LangGraph
# (get_graph() + NODE_META sidecar) — new agents appear auto
status_page.py # read-only DATA LAYER for /api/state (build_snapshot /
# snapshot_to_dict); HTML rendering retired in the makeover
billing.py # §3.1 claude_invoke billing-mode seam
ci_gate.py # §3.3.2 pure-code authenticated-Checks PASS/FAIL gate
task_model.py / state_store.py
db/{schema.py,schema.sql} # SQLite ledger DDL + BEGIN IMMEDIATE compare-and-set
db/transitions.py # task_transitions recorder (per-task pipeline history,
# idempotent + fail-soft; written by instrumented graph nodes)
ledger.py / responder.py / resume_worker.py / deadline_timer.py / recovery.py
operator_cli.py
nodes/ # pipeline stages + their model bindings
@ -64,9 +72,11 @@ agent-team/
slack_adapter.py + slack_live.py + slack_listener.py # Block Kit + Socket Mode + /new-task
github_adapter.py + github_live.py + github_intake.py # issue-comment + issue intake
claude_code_adapter.py + claude_code_live.py # file-drop responder
web/ # React/Vite/TypeScript status-dashboard SPA (React Flow map,
# task list, click-through task history); built to web/dist
scripts/ # deploy-r720-ws-rollout.sh — attended WS0–WS5 UPDATE of the box
ci/ # §3.3.2 split-job CI apply/verify workflow (DEPLOY-GATED)
systemd/ # agent-team-coordinator.service (not installed)
systemd/ # agent-team-coordinator.service + agent-team-status.service
DEPLOY-R720.md # provisioning runbook (snapshot-first, rsync, tokens, demo)
tests/ # pytest, one module per source module + sim harness
```
@ -98,6 +108,37 @@ snake_case (`agent_team/`), per the engineering handbook.
GPT-4.1 (review) and DeepSeek (builders) route through the local orchestrator
`run.py`. Switching Claude billing is a config flip.
## Status dashboard (WebUI)
A read-only LAN dashboard (FastAPI, `0.0.0.0:8770`, no auth, `mode=ro` ledger
opens) for watching the pipeline. Served by `agent_team.dashboard` (systemd unit
`agent-team-status.service`):
- **Live pipeline map** — a React Flow graph **auto-laid-out from the real
LangGraph** (`topology.py` introspects `compiled.get_graph()` + a `NODE_META`
display sidecar). Adding an agent node in `graph.py` makes it appear on the map
with no manual coordinates; nodes group into **trees** (processes) branching off
`intake`. Node color = live state; loop-back edges (review→plan, verify→build)
render dashed.
- **Click-through task history** — selecting a task opens a timeline of its journey
through each node (entry/exit timestamps, per-node duration, per-stage cost, Q&A,
verdicts, plan), backed by the `task_transitions` ledger (schema v3) written by
the coordinator's instrumented graph nodes (`db/transitions.py`, fail-soft).
Endpoints: `GET /api/state` (live overview + per-node state), `GET /api/topology`
(map nodes/edges/trees), `GET /api/task/{thread_id}` (one task's history;
`thread_id` is validated `^[A-Za-z0-9_-]{1,64}$`). The legacy stdlib HTML page was
retired; `status_page.py` remains as the `/api/state` data layer.
**Build (on the Mac — the box Node is too old for Vite 5+):**
```
cd agent-team/web
npm ci && npm run build # -> web/dist (gitignored), rsynced to the VM
npm test # Vitest + React Testing Library
```
In dev, `npm run dev` proxies `/api` to a locally running dashboard
(`AGENT_TEAM_DASH`, default `http://127.0.0.1:8770`).
## WS0–WS5 rollout glossary
The "WS-rollout" (workstreams 0–5) layered HTTP/integration surfaces onto the

View file

@ -60,6 +60,7 @@ from typing import TYPE_CHECKING, Any, Callable
from agent_team import graph as graph_mod
from agent_team import responder as responder_mod
from agent_team.db.schema import connect, init_db, supersede_question
from agent_team.db.transitions import TransitionRecorder
from agent_team.resume_worker import ResumeOutcome, ResumeResult, ResumeWorker
from agent_team.transport.base import Transport
from agent_team.transport.slack_adapter import SlackTransport
@ -604,8 +605,16 @@ class Coordinator:
if self._dispatch_node_wiring is not None:
dispatch_node_callable = self._dispatch_node_wiring()
# Per-task transition history (dashboard drill-down): instrument every
# graph node to append a task_transitions row on entry. Fail-soft — a
# ledger write never breaks the pipeline — so this is always wired in
# the live coordinator. Tests that build the graph directly default to
# transition_recorder=None (no instrumentation).
transition_recorder = TransitionRecorder(self._db_path)
self._graph = graph_mod.build_graph(
checkpointer,
transition_recorder=transition_recorder,
live_clarify_node=clarify_node,
live_plan_node=plan_node,
review_node=review_node,

View file

@ -0,0 +1,305 @@
"""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 only the exception TYPE, never str(exc) — a SQLite message can
# carry the DB path / table names; keep the LAN surface tight.
return {
"ok": False,
"error": f"read failed: {type(exc).__name__}",
"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
# Type only, never str(exc) — keep internal LangGraph/import detail
# off the unauthenticated LAN surface (SEC-DASH-003).
payload = {
"trees": [],
"nodes": [],
"edges": [],
"error": type(exc).__name__,
}
return JSONResponse(payload, headers={"Cache-Control": "no-store"})
github-advanced-security[bot] commented 2026-06-23 21:21:11 +00:00 (Migrated from github.com)
Review

CodeQL / Information exposure through an exception

Stack trace information flows to this location and may be exposed to an external user.

Show more details

## CodeQL / Information exposure through an exception [Stack trace information](1) flows to this location and may be exposed to an external user. [Show more details](https://github.com/Sea-Haven-Industries/orchestrator/security/code-scanning/6)
@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)

View file

@ -9,16 +9,23 @@ from agent_team.db.schema import (
BUDGET_LEDGER_DDL,
PENDING_QUESTIONS_DDL,
SCHEMA_VERSION,
TASK_TRANSITIONS_DDL,
assert_task_transitions_ready,
connect,
init_db,
migrate,
)
from agent_team.db.transitions import TransitionRecorder, read_transitions
__all__ = [
"BUDGET_LEDGER_DDL",
"PENDING_QUESTIONS_DDL",
"SCHEMA_VERSION",
"TASK_TRANSITIONS_DDL",
"TransitionRecorder",
"assert_task_transitions_ready",
"connect",
"init_db",
"migrate",
"read_transitions",
]

View file

@ -32,10 +32,13 @@ __all__ = [
"PENDING_QUESTIONS_INDEXES_DDL",
"BUDGET_LEDGER_INDEXES_DDL",
"INGESTED_ISSUES_DDL",
"TASK_TRANSITIONS_DDL",
"TASK_TRANSITIONS_INDEXES_DDL",
"SCHEMA_META_DDL",
"SCHEMA_VERSION",
"QUESTION_STATES",
"answer_question",
"assert_task_transitions_ready",
"connect",
"delete_issue_ingested",
"expire_question",
@ -49,7 +52,7 @@ __all__ = [
]
# Bump when the DDL below changes; migrate() steps a connection forward.
SCHEMA_VERSION: int = 2
SCHEMA_VERSION: int = 3
# Default SQLite busy timeout (ms) so concurrent writers wait for the write
# lock rather than failing immediately.
@ -127,6 +130,35 @@ CREATE TABLE IF NOT EXISTS ingested_issues (
)
""".strip()
# task_transitions: per-task pipeline history (schema v3). One row per node a
# task enters, recording the move from_phase -> to_phase with entry/exit
# timestamps and the task status at entry. The dashboard's /api/task drill-down
# reads this to render a task's journey through the pipeline (timestamps,
# per-node duration); per-stage cost is joined from budget_ledger at read time,
# not duplicated here. Written by the coordinator's instrumented graph nodes
# (db/transitions.py); read READ-ONLY by the status dashboard. ``exited_at`` is
# filled when the next transition lands OR when the task reaches a terminal
# status (the recorder's close_terminal). A single OPEN row per thread is the
# invariant the recorder maintains (close-open-before-insert), so a crash/resume
# replay cannot leave orphaned open rows.
TASK_TRANSITIONS_DDL: str = """
CREATE TABLE IF NOT EXISTS task_transitions (
transition_id INTEGER PRIMARY KEY AUTOINCREMENT,
thread_id TEXT NOT NULL,
from_phase TEXT,
to_phase TEXT NOT NULL,
entered_at TEXT NOT NULL,
exited_at TEXT,
status TEXT,
note TEXT
)
""".strip()
TASK_TRANSITIONS_INDEXES_DDL: str = """
CREATE INDEX IF NOT EXISTS idx_task_transitions_thread
ON task_transitions (thread_id, entered_at);
""".strip()
SCHEMA_META_DDL: str = """
CREATE TABLE IF NOT EXISTS schema_meta (
id INTEGER PRIMARY KEY CHECK (id = 1),
@ -203,6 +235,9 @@ def init_db(db_path: Path) -> None:
for stmt in _split_statements(BUDGET_LEDGER_INDEXES_DDL):
conn.execute(stmt)
conn.execute(INGESTED_ISSUES_DDL)
conn.execute(TASK_TRANSITIONS_DDL)
for stmt in _split_statements(TASK_TRANSITIONS_INDEXES_DDL):
conn.execute(stmt)
# Record the schema version (single-row table). DO NOTHING leaves an
# existing row's version untouched (an already-stamped DB just gains any
# IF-NOT-EXISTS tables above); migrate() is what steps the version stamp
@ -244,12 +279,28 @@ def migrate(conn: sqlite3.Connection) -> None:
conn.execute(INGESTED_ISSUES_DDL)
current = 2
# Future steps go here: `if current < 3: ...; current = 3`.
if current < 3:
# v3: per-task pipeline transition history (dashboard drill-down).
conn.execute(TASK_TRANSITIONS_DDL)
for stmt in _split_statements(TASK_TRANSITIONS_INDEXES_DDL):
conn.execute(stmt)
current = 3
# Future steps go here: `if current < 4: ...; current = 4`.
# Applied UNCONDITIONALLY (idempotent IF NOT EXISTS) so an already-stamped DB
# — which skips the version blocks above — still gains the v2 table without a
# restamp. Safe on existing data: a fresh empty table only.
# — which skips the version blocks above — still gains these tables without a
# restamp. Safe on existing data: fresh empty tables / indexes only.
#
# ORDERING CONSTRAINT: a future migration that ALTERs one of these tables
# (e.g. `ALTER TABLE task_transitions ADD COLUMN ...`) MUST run in its own
# `if current < N` block placed ABOVE this tail — the unconditional CREATE
# ... IF NOT EXISTS here no-ops on an existing table and will NOT apply an
# alter. This tail is only for first-time creation on an already-stamped DB.
conn.execute(INGESTED_ISSUES_DDL)
conn.execute(TASK_TRANSITIONS_DDL)
for stmt in _split_statements(TASK_TRANSITIONS_INDEXES_DDL):
conn.execute(stmt)
# Defense-in-depth index, applied UNCONDITIONALLY (idempotent IF NOT EXISTS)
# so an already-stamped v1 DB — which skips the `current < 1` block above —
@ -271,6 +322,51 @@ def migrate(conn: sqlite3.Connection) -> None:
)
# Columns task_transitions MUST expose for the dashboard drill-down to work.
# assert_task_transitions_ready checks these so a botched migration is caught at
# startup (loud) rather than surfacing as a half-broken /api/task at read time.
_TASK_TRANSITIONS_COLUMNS: frozenset[str] = frozenset(
{
"transition_id",
"thread_id",
"from_phase",
"to_phase",
"entered_at",
"exited_at",
"status",
"note",
}
)
def assert_task_transitions_ready(db_path: Path) -> None:
"""Assert the v3 ``task_transitions`` table exists with the expected columns.
The deploy convention (and BLOCK-2 in the WebUI-makeover plan) requires the
migration to be verified *before* declaring a deploy good — a missing or
malformed table must fail loudly at startup, not after a crash-loop or as a
silently-broken drill-down. Raises :class:`RuntimeError` if the table is
absent or any expected column is missing; returns ``None`` on success.
"""
conn = connect(db_path)
try:
rows = conn.execute("PRAGMA table_info(task_transitions)").fetchall()
finally:
conn.close()
if not rows:
raise RuntimeError(
f"task_transitions table missing in {db_path} after migration "
"(schema v3) — aborting; restore the ledger backup and re-deploy."
)
present = {row["name"] for row in rows}
missing = _TASK_TRANSITIONS_COLUMNS - present
if missing:
raise RuntimeError(
f"task_transitions in {db_path} is malformed — missing columns "
f"{sorted(missing)}; aborting, restore the ledger backup."
)
def issue_already_ingested(
conn: sqlite3.Connection, *, source: str, issue_id: str
) -> bool:

View file

@ -0,0 +1,206 @@
"""Per-task pipeline transition recorder (schema v3, ``task_transitions``).
The coordinator wraps each LangGraph node (``graph._instrument``) so that, as a
task enters a node, one row is appended to ``task_transitions`` capturing the
move ``from_phase -> to_phase`` with an ``entered_at`` stamp and the task status
at entry. The dashboard's ``/api/task/{thread_id}`` drill-down reads these rows
to render a task's journey through the pipeline; per-stage cost is joined from
``budget_ledger`` at read time and is NOT duplicated here.
Two design properties matter (both exercised by tests):
* **Fail-soft.** Every write is wrapped so a ledger problem (locked DB, missing
table, disk error) is logged and swallowed — instrumentation must NEVER break
the live pipeline. A dropped transition row degrades the dashboard, nothing
more.
* **Idempotent under replay.** LangGraph re-executes a node from its start on
resume (e.g. the clarifier replays after the human gate). So a node's wrapper
may call :meth:`TransitionRecorder.record_entry` more than once for the same
``(thread_id, to_phase)``. The recorder keeps a single OPEN row per thread:
re-entering the *same* node while it is already the open row is a no-op;
entering a *different* node closes the previous open row (filling its
``exited_at``) before inserting the new one. A crash/resume therefore cannot
leave orphaned open rows or double-count a replayed node.
"""
from __future__ import annotations
import logging
import sqlite3
from datetime import datetime, timezone
from pathlib import Path
from agent_team.db.schema import connect
__all__ = ["TransitionRecorder", "read_transitions"]
_log = logging.getLogger(__name__)
def _utc_now_iso() -> str:
"""Current UTC time as an ISO-8601 string (matches the other ledgers)."""
return datetime.now(timezone.utc).isoformat()
def read_transitions(
conn: sqlite3.Connection, thread_id: str
) -> list[dict[str, object]]:
"""Return ``thread_id``'s transition rows in entry order (oldest first).
Read-only and parameterized; callers pass their own connection (the
dashboard uses a strictly read-only one). Returns ``[]`` if the table does
not exist yet (fresh DB) rather than raising, so a status view never crashes.
"""
try:
rows = conn.execute(
"SELECT transition_id, thread_id, from_phase, to_phase, "
"entered_at, exited_at, status, note "
"FROM task_transitions WHERE thread_id = ? "
"ORDER BY entered_at ASC, transition_id ASC",
(thread_id,),
).fetchall()
except sqlite3.Error:
# Missing table (fresh DB) or a corrupt/unreadable ledger — a status
# view must never crash on a read.
return []
return [dict(row) for row in rows]
class TransitionRecorder:
"""Fail-soft writer for ``task_transitions`` (one OPEN row per thread).
Constructed with the ledger DB **path** (not a connection): the coordinator
owns the writable side, and each call opens a short-lived WAL connection via
:func:`agent_team.db.schema.connect`, mirroring the compare-and-set helpers'
connection discipline. An in-memory path (``":memory:"``) is supported for
tests by reusing a single retained connection (a fresh ``:memory:`` connect
would see an empty database).
"""
def __init__(self, db_path: Path | str) -> None:
self._db_path = Path(db_path)
self._in_memory = str(db_path) == ":memory:"
# In-memory DBs are per-connection; retain one so writes accumulate.
self._mem_conn: sqlite3.Connection | None = (
connect(self._db_path) if self._in_memory else None
)
def _connect(self) -> sqlite3.Connection:
if self._mem_conn is not None:
return self._mem_conn
return connect(self._db_path)
def _close(self, conn: sqlite3.Connection) -> None:
# Never close the retained in-memory connection.
if conn is not self._mem_conn:
conn.close()
def close(self) -> None:
"""Close the retained in-memory connection, if any (test cleanup).
File-backed recorders open/close per call and hold nothing, so this is a
no-op for them; the in-memory test path retains one connection that this
releases so it does not leak when the recorder is discarded.
"""
if self._mem_conn is not None:
self._mem_conn.close()
self._mem_conn = None
def record_entry(
self,
*,
thread_id: str,
to_phase: str,
status: str | None = None,
note: str | None = None,
) -> None:
"""Append an OPEN transition for entering ``to_phase`` (idempotent).
If the thread's latest open row is already this ``to_phase`` (a resume
replay of the same node), this is a no-op. If the latest open row is a
*different* node, it is closed (``exited_at`` filled) before the new open
row is inserted, carrying that node's ``to_phase`` forward as the new
row's ``from_phase``. Fail-soft: any error is logged and swallowed.
"""
if not thread_id or not to_phase:
return
conn = None
try:
conn = self._connect()
now = _utc_now_iso()
open_row = conn.execute(
"SELECT transition_id, to_phase FROM task_transitions "
"WHERE thread_id = ? AND exited_at IS NULL "
"ORDER BY transition_id DESC LIMIT 1",
(thread_id,),
).fetchone()
from_phase: str | None = None
if open_row is not None:
if open_row["to_phase"] == to_phase:
# Same node re-entered (resume replay) — no double row.
return
# Different node: close the previous open row.
conn.execute(
"UPDATE task_transitions SET exited_at = ? WHERE transition_id = ?",
(now, open_row["transition_id"]),
)
from_phase = open_row["to_phase"]
conn.execute(
"INSERT INTO task_transitions "
"(thread_id, from_phase, to_phase, entered_at, exited_at, "
"status, note) VALUES (?, ?, ?, ?, NULL, ?, ?)",
(thread_id, from_phase, to_phase, now, status, note),
)
except Exception as exc: # fail-soft: never break the pipeline
_log.warning("task_transitions record_entry failed: %s", exc)
finally:
if conn is not None:
self._close(conn)
def close_terminal(
self,
*,
thread_id: str,
status: str | None = None,
note: str | None = None,
) -> None:
"""Close the thread's open transition on a terminal status (N2).
The last node a task visits never gets an ``exited_at`` from a *next*
transition, so when a node returns a terminal status (DONE/PARKED/FAILED)
the wrapper calls this to stamp ``exited_at`` (and optionally update the
row's ``status``/``note``). No-op if there is no open row. Fail-soft.
"""
if not thread_id:
return
conn = None
try:
conn = self._connect()
now = _utc_now_iso()
open_row = conn.execute(
"SELECT transition_id FROM task_transitions "
"WHERE thread_id = ? AND exited_at IS NULL "
"ORDER BY transition_id DESC LIMIT 1",
(thread_id,),
).fetchone()
if open_row is None:
return
if status is not None:
conn.execute(
"UPDATE task_transitions SET exited_at = ?, status = ?, "
"note = COALESCE(?, note) WHERE transition_id = ?",
(now, status, note, open_row["transition_id"]),
)
else:
conn.execute(
"UPDATE task_transitions SET exited_at = ?, "
"note = COALESCE(?, note) WHERE transition_id = ?",
(now, note, open_row["transition_id"]),
)
except Exception as exc: # fail-soft
_log.warning("task_transitions close_terminal failed: %s", exc)
finally:
if conn is not None:
self._close(conn)

View file

@ -38,6 +38,8 @@ shapes the interrupt payload that drives it.
from __future__ import annotations
import functools
import inspect
import uuid
from collections.abc import Callable
from contextlib import contextmanager
@ -308,12 +310,72 @@ def plan_phase(state: PipelineState) -> dict[str, Any]:
}
# --- Node instrumentation (per-task transition history). --------------------
# Terminal status VALUES (TaskStatus.value) that close a task's open transition
# row — the last node never gets an exited_at from a *next* transition (plan §N2).
_TERMINAL_STATUS_VALUES: frozenset[str] = frozenset(
{TaskStatus.DONE.value, TaskStatus.PARKED.value, TaskStatus.FAILED.value}
)
def _instrument(
name: str,
fn: Callable[..., Any],
recorder: Any | None,
) -> Callable[..., Any]:
"""Wrap a graph node so entering it records a ``task_transitions`` row.
Returns ``fn`` unchanged when ``recorder`` is ``None`` (the default — current
tests and the uninstrumented graph are untouched). Otherwise returns a
signature-preserving wrapper that, on entry, calls ``recorder.record_entry``
(idempotent under LangGraph's resume replay) and, when the node returns a
terminal status, calls ``recorder.close_terminal`` to stamp ``exited_at``.
**Signature preservation (plan §B4).** LangGraph's ``add_node`` inspects the
callable's signature to decide whether to inject a ``RunnableConfig`` second
argument (there is a prior fixed bug of this exact class). ``functools.wraps``
sets ``__wrapped__`` (which ``inspect.signature`` follows) and we ALSO set
``__signature__`` explicitly to ``fn``'s, so LangGraph sees ``fn``'s real
arity and passes exactly the arguments ``fn`` expects; the wrapper forwards
them verbatim via ``*args, **kwargs``. Recording is best-effort: the recorder
is itself fail-soft, and the node call is never gated on it.
"""
if recorder is None:
return fn
@functools.wraps(fn)
def wrapped(state: Any, *args: Any, **kwargs: Any) -> Any:
thread_id = ""
status: Any = None
if isinstance(state, dict):
thread_id = state.get("thread_id", "") or ""
status = state.get("status")
recorder.record_entry(thread_id=thread_id, to_phase=name, status=status)
result = fn(state, *args, **kwargs)
if isinstance(result, dict):
# Nodes store status as TaskStatus.value strings; coerce defensively
# so an enum member (should one slip through) still closes the row.
new_status = result.get("status")
status_value = getattr(new_status, "value", new_status)
if status_value in _TERMINAL_STATUS_VALUES:
recorder.close_terminal(thread_id=thread_id, status=status_value)
return result
# Belt-and-suspenders for B4: present fn's exact signature to LangGraph.
try:
wrapped.__signature__ = inspect.signature(fn) # type: ignore[attr-defined]
except (TypeError, ValueError): # pragma: no cover - exotic callables
pass
return wrapped
# --- Graph assembly. --------------------------------------------------------
def build_graph(
checkpointer: BaseCheckpointSaver | None = None,
*,
transition_recorder: Any | None = None,
live_clarify_node: Callable[[PipelineState], PipelineState] | None = None,
live_plan_node: Callable[[PipelineState], PipelineState] | None = None,
review_node: Callable[[PipelineState], PipelineState] | None = None,
@ -416,9 +478,9 @@ def build_graph(
)
builder: StateGraph = StateGraph(PipelineState)
builder.add_node(INTAKE, intake_node)
builder.add_node(CLARIFY, clarify)
builder.add_node(PLAN, plan)
builder.add_node(INTAKE, _instrument(INTAKE, intake_node, transition_recorder))
builder.add_node(CLARIFY, _instrument(CLARIFY, clarify, transition_recorder))
builder.add_node(PLAN, _instrument(PLAN, plan, transition_recorder))
builder.add_edge(START, INTAKE)
builder.add_edge(INTAKE, CLARIFY)
@ -429,7 +491,7 @@ def build_graph(
builder.add_edge(PLAN, END)
else:
# P2/P3: plan -> review -> {loop-back to plan | build | END}.
builder.add_node(REVIEW, review_node)
builder.add_node(REVIEW, _instrument(REVIEW, review_node, transition_recorder))
builder.add_edge(PLAN, REVIEW)
if build_verify is None:
@ -446,8 +508,12 @@ def build_graph(
# END (escalation)}. The subgraph nodes + router are injected (the
# ``build_verify`` tuple) so this module imports no P3 code.
build_node, verify_node, route_after_verify = build_verify
builder.add_node(BUILD_NODE, build_node)
builder.add_node(VERIFY_NODE, verify_node)
builder.add_node(
BUILD_NODE, _instrument(BUILD_NODE, build_node, transition_recorder)
)
builder.add_node(
VERIFY_NODE, _instrument(VERIFY_NODE, verify_node, transition_recorder)
)
builder.add_conditional_edges(
REVIEW,
@ -464,7 +530,10 @@ def build_graph(
)
else:
# P3+: repoint APPROVED_ROUTE at the dispatch node, then END.
builder.add_node(DISPATCH_NODE, dispatch_node)
builder.add_node(
DISPATCH_NODE,
_instrument(DISPATCH_NODE, dispatch_node, transition_recorder),
)
builder.add_conditional_edges(
VERIFY_NODE,
route_after_verify,

View file

@ -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"}
@ -313,19 +298,23 @@ def build_snapshot(db_path: Path | str | None = None) -> Snapshot:
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.
# Locked / busy / unreadable — never crash, just say so. Return only the
# exception TYPE, never str(exc): a SQLite message can carry the DB path /
# table names, and this payload is served unauthenticated on the LAN
# (matches the task_detail hardening). [sh-security-review SEC-DASH-001]
return Snapshot(
generated_at=generated_at,
db_path=str(resolved),
ok=False,
error=f"could not read ledger (busy or unreadable): {exc}",
error=f"could not read ledger (busy or unreadable): {type(exc).__name__}",
)
except Exception as exc: # defensive: any unexpected read failure
# Type only, never str(exc) — see SEC-DASH-001 note above.
return Snapshot(
generated_at=generated_at,
db_path=str(resolved),
ok=False,
error=f"unexpected error reading ledger: {exc}",
error=f"unexpected error reading ledger: {type(exc).__name__}",
)
finally:
if conn is not None:
@ -587,711 +576,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 "&mdash;"
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 "&mdash;"
)
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 '&mdash;')}</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 &mdash; 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 &mdash; 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> &middot; live "
f"(polls /api/state every {_POLL_INTERVAL_MS // 1000}s) &middot; "
f'<span class="mono">{_esc(snapshot.db_path)}</span> &middot; '
"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 &middot; no mutating endpoints "
"&middot; unauthenticated &mdash; 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.

View 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,
}

View file

@ -1,15 +1,20 @@
# agent-team-status.service - R720 LAN-only READ-ONLY status dashboard (sh-secrev VM, user adam).
#
# A tiny stdlib http.server that renders the agent-team coordinator's queue
# (tasks/phases, who is waiting on the human gate, active/parked counts, recent
# budget spend) as a 10s-auto-refresh HTML page. It opens the SQLite ledger
# READ-ONLY (mode=ro) and never writes; it has no mutating endpoints and no auth.
# Serves the React/Vite single-page dashboard (web/dist) + its read-only JSON API
# (/api/state, /api/topology, /api/task/{id}) via FastAPI/uvicorn. The map renders
# the live pipeline (auto-laid from the real LangGraph), and a task can be clicked
# to see its history through each node. It opens the SQLite ledger READ-ONLY
# (mode=ro) and never writes; it has no mutating endpoints and no auth.
#
# NETWORK POSTURE: binds AGENT_TEAM_STATUS_HOST (default 0.0.0.0) on port
# AGENT_TEAM_STATUS_PORT (default 8770). The sh-secrev VM (10.10.60.120, VLAN 60)
# has NO public NIC and sits behind the UniFi firewall, so 0.0.0.0 reaches the
# LAN/VPN only. Task descriptions may be sensitive -> keep this LAN/VPN-only,
# never expose to the public internet.
# NETWORK POSTURE: binds 0.0.0.0 on port 8770. The sh-secrev VM (10.10.60.120,
# VLAN 60) has NO public NIC and sits behind the UniFi firewall, so 0.0.0.0
# reaches the LAN/VPN only. Task descriptions may be sensitive -> keep this
# LAN/VPN-only, never expose to the public internet.
#
# BUILD: the SPA is built on the Mac (`cd web && npm ci && npm run build`) and the
# resulting web/dist is rsynced to the VM (the box Node is too old for a Vite 5+
# build). uvicorn serves whatever web/dist is present; if absent, the JSON API
# still works and the SPA 404s until dist is deployed.
#
# Install (on the VM, as root):
# sudo cp agent-team-status.service /etc/systemd/system/
@ -35,9 +40,10 @@ WorkingDirectory=/home/adam/orchestrator/agent-team
# as the coordinator lets AGENT_TEAM_DB / AGENT_TEAM_STATUS_* overrides live in
# one place if set there.
EnvironmentFile=-/home/adam/secrev.env
# Use the agent-team venv interpreter (where langgraph + the checkpoint dep are
# installed), NOT the bare system python3 that systemd's PATH would resolve.
ExecStart=/home/adam/orchestrator/agent-team/.venv/bin/python -c "from agent_team.status_page import serve; serve()"
# Use the agent-team venv interpreter (where langgraph + fastapi/uvicorn + the
# checkpoint dep are installed), NOT the bare system python3 that systemd's PATH
# would resolve.
ExecStart=/home/adam/orchestrator/agent-team/.venv/bin/python -c "from agent_team.dashboard import serve; serve()"
Restart=on-failure
RestartSec=5
# Hardening - mirrors agent-team-coordinator.service, but this unit only READS

View file

@ -11,3 +11,9 @@ from pathlib import Path
_PROJECT_ROOT = Path(__file__).resolve().parents[1]
if str(_PROJECT_ROOT) not in sys.path:
sys.path.insert(0, str(_PROJECT_ROOT))
# Also expose the tests/ dir so a test module can import shared helpers from a
# sibling test module by bare name (e.g. ``from test_status_page import ...``).
_TESTS_DIR = Path(__file__).resolve().parent
if str(_TESTS_DIR) not in sys.path:
sys.path.insert(0, str(_TESTS_DIR))

View file

@ -0,0 +1,118 @@
"""Tests for the read-only dashboard FastAPI app (``agent_team.dashboard``).
Exercises the three JSON endpoints against a seeded temp ledger via FastAPI's
TestClient: /api/state (contract preserved + per-node live state), /api/topology,
and /api/task/{id} (timeline assembly, thread_id validation, partial fallback,
cost join). Guarded so the suite still runs where FastAPI is absent.
"""
from __future__ import annotations
import importlib.util
from pathlib import Path
import pytest
from agent_team.db import TransitionRecorder, connect, init_db
from agent_team.dashboard import task_detail
# Reuse the seeded-ledger helper from the status-page tests (sibling module,
# importable by bare name under pytest's default prepend import mode).
from test_status_page import seed_ledger
_HAS_FASTAPI = importlib.util.find_spec("fastapi") is not None
pytestmark = pytest.mark.skipif(_HAS_FASTAPI is False, reason="fastapi not installed")
def _client(db: Path):
from fastapi.testclient import TestClient
from agent_team.dashboard import make_dashboard_app
return TestClient(make_dashboard_app(db_path=db))
def _seed_budget(db: Path, thread_id: str, stage: str, usd: float) -> None:
conn = connect(db)
try:
conn.execute(
"INSERT INTO budget_ledger (thread_id, stage, model, billing_mode, "
"usd_cost, recorded_at, day_bucket) VALUES (?,?,?,?,?,?,?)",
(
thread_id,
stage,
"claude",
"subscription",
usd,
"2026-06-23T00:00:00+00:00",
"2026-06-23",
),
)
finally:
conn.close()
def test_api_state_contract_and_nodes(tmp_path: Path) -> None:
db = tmp_path / "agent_team.sqlite"
seed_ledger(db)
r = _client(db).get("/api/state")
assert r.status_code == 200
payload = r.json()
# Existing contract preserved.
assert payload["ok"] is True
assert {"stages", "tasks", "summary", "budget"} <= set(payload)
# New per-node live state, keyed by topology node id.
assert "nodes" in payload
# thread-waiting is at clarify and awaiting human; thread-active at plan.
assert payload["nodes"]["clarify"]["state"] == "awaiting_human"
assert payload["nodes"]["plan"]["state"] == "active"
def test_api_topology(tmp_path: Path) -> None:
db = tmp_path / "agent_team.sqlite"
init_db(db)
payload = _client(db).get("/api/topology").json()
ids = {n["id"] for n in payload["nodes"]}
assert {"intake", "clarify", "plan", "review"} <= ids
assert payload["edges"]
def test_api_task_timeline_and_cost(tmp_path: Path) -> None:
db = tmp_path / "agent_team.sqlite"
seed_ledger(db)
rec = TransitionRecorder(db)
rec.record_entry(thread_id="thread-active", to_phase="intake", status="active")
rec.record_entry(thread_id="thread-active", to_phase="plan", status="active")
_seed_budget(db, "thread-active", "plan", 0.5)
payload = _client(db).get("/api/task/thread-active").json()
assert payload["ok"] is True
assert payload["partial"] is False
steps = {s["to_phase"]: s for s in payload["timeline"]}
assert "intake" in steps and "plan" in steps
assert steps["plan"]["cost_usd"] == pytest.approx(0.5)
assert payload["total_usd"] == pytest.approx(0.5)
def test_api_task_partial_when_no_transitions(tmp_path: Path) -> None:
db = tmp_path / "agent_team.sqlite"
seed_ledger(db) # checkpoints exist, but no task_transitions rows
payload = _client(db).get("/api/task/thread-active").json()
assert payload["ok"] is True
assert payload["partial"] is True
# Best-effort single step reconstructed from the checkpoint phase.
assert payload["timeline"]
assert payload["timeline"][0]["note"]
def test_api_task_rejects_bad_thread_id(tmp_path: Path) -> None:
db = tmp_path / "agent_team.sqlite"
init_db(db)
cli = _client(db)
assert cli.get("/api/task/has space").status_code == 400
assert cli.get("/api/task/" + "x" * 65).status_code == 400
def test_task_detail_missing_ledger_is_failsafe(tmp_path: Path) -> None:
out = task_detail(tmp_path / "nope.sqlite", "abc")
assert out["ok"] is False

View file

@ -1,10 +1,12 @@
"""Tests for the LAN-only read-only status dashboard (``agent_team.status_page``).
"""Tests for the status dashboard DATA LAYER (``agent_team.status_page``).
No live socket is bound: :func:`render_html` is exercised against fabricated
:class:`Snapshot` objects, and the read-only reader (:func:`build_snapshot`) is
exercised against a temp SQLite ledger seeded with a couple of rows (a real
``SqliteSaver`` checkpoint plus ``pending_questions`` / ``budget_ledger`` rows).
``conftest.py`` already puts ``agent-team/`` on ``sys.path``.
After the WebUI makeover, ``status_page`` no longer renders HTML — the dashboard
is a React/Vite SPA served by ``agent_team.dashboard`` (see ``test_dashboard.py``
for the FastAPI endpoints). What remains here is the read-only reader
(:func:`build_snapshot`) and the ``/api/state`` payload builder
(:func:`snapshot_to_dict`), exercised against fabricated :class:`Snapshot`
objects and a temp SQLite ledger seeded with real rows. ``conftest.py`` puts
``agent-team/`` on ``sys.path``.
"""
from __future__ import annotations
@ -21,12 +23,9 @@ from agent_team.status_page import (
Snapshot,
TaskView,
build_snapshot,
render_html,
snapshot_to_dict,
)
# --- pure render_html tests (no DB, no socket). ------------------------------
def _sample_snapshot() -> Snapshot:
return Snapshot(
@ -74,85 +73,16 @@ def _sample_snapshot() -> Snapshot:
)
def test_render_html_returns_complete_document() -> None:
html_out = render_html(_sample_snapshot())
assert html_out.startswith("<!doctype html>")
assert html_out.rstrip().endswith("</html>")
# Auto-refresh + offline (no external CDN / http references).
assert '<meta http-equiv="refresh" content="10">' in html_out
assert "http://" not in html_out and "https://" not in html_out
assert "<style>" in html_out # inline CSS only
def test_render_html_shows_counts_and_waiting() -> None:
snap = _sample_snapshot()
html_out = render_html(snap)
assert snap.active_count == 2 # active + waiting_human
assert snap.parked_count == 1
assert len(snap.waiting_tasks) == 1
# The waiting task and its phase are surfaced.
assert "abc123de" in html_out
assert "remediate CVE-2026-0001 in payments-svc" in html_out
assert "Waiting on the human gate" in html_out
# Budget panel present with a total.
assert "Budget" in html_out
assert "claude-opus-4" in html_out
def test_render_html_escapes_task_descriptions() -> None:
snap = Snapshot(
generated_at="now",
db_path="/x.sqlite",
ok=True,
tasks=[
TaskView(
thread_id="x",
short_id="x",
task="<script>alert('xss')</script>",
current_phase="intake",
status="active",
)
],
)
html_out = render_html(snap)
assert "<script>alert" not in html_out
assert "&lt;script&gt;" in html_out
def test_render_html_no_data_page_on_not_ok() -> None:
snap = Snapshot(
generated_at="now",
db_path="/missing.sqlite",
ok=False,
error="ledger not found",
)
html_out = render_html(snap)
assert "No data" in html_out
assert "ledger not found" in html_out
assert html_out.rstrip().endswith("</html>")
def test_render_html_omits_budget_when_unavailable() -> None:
snap = Snapshot(
generated_at="now",
db_path="/x.sqlite",
ok=True,
tasks=[],
budget_available=False,
)
html_out = render_html(snap)
assert "Budget" not in html_out
# --- read-only reader tests against a temp seeded SQLite ledger. -------------
def _seed_db(db_path: Path) -> None:
def seed_ledger(db_path: Path) -> None:
"""Create the agent-team schema and seed a couple of rows.
Seeds one pending question (open) and one budget row via the foundation
Seeds one open pending question and one budget row via the foundation
schema, and one LangGraph checkpoint per thread via the real ``SqliteSaver``
so the reader's checkpoint path is exercised faithfully.
so the reader's checkpoint path is exercised faithfully. Imported by
``test_dashboard.py`` as well.
"""
from agent_team.db.schema import init_db
from langgraph.checkpoint.base import empty_checkpoint
@ -183,7 +113,6 @@ def _seed_db(db_path: Path) -> None:
)
conn.commit()
# Seed checkpoints for two threads via the real saver.
with closing(sqlite3.connect(str(db_path), check_same_thread=False)) as conn:
saver = SqliteSaver(conn)
for tid, values in (
@ -216,13 +145,12 @@ def _seed_db(db_path: Path) -> None:
def test_build_snapshot_reads_seeded_ledger(tmp_path: Path) -> None:
db = tmp_path / "agent_team.sqlite"
_seed_db(db)
seed_ledger(db)
snap = build_snapshot(db)
assert snap.ok is True
assert snap.error is None
# Two checkpointed threads enumerated.
by_id = {t.thread_id: t for t in snap.tasks}
assert set(by_id) == {"thread-waiting", "thread-active"}
@ -230,35 +158,26 @@ def test_build_snapshot_reads_seeded_ledger(tmp_path: Path) -> None:
assert waiting.task == "remediate CVE"
assert waiting.current_phase == "clarify"
assert waiting.status == "waiting_human"
assert waiting.waiting is True # has an open pending_question
assert waiting.waiting is True
assert waiting.waiting_since == "2026-06-23T11:50:00+00:00"
active = by_id["thread-active"]
assert active.waiting is False
assert active.status == "active"
# Queue summary + question counts.
assert snap.active_count == 2 # active + waiting_human
assert snap.parked_count == 0
assert snap.question_counts.get("open") == 1
# Budget read.
assert snap.budget_available is True
assert snap.spend_total_usd == pytest.approx(0.25)
assert snap.recent_spend[0]["model"] == "claude-opus-4"
# The whole thing renders without error.
html_out = render_html(snap)
assert "remediate CVE" in html_out
def test_build_snapshot_missing_db_is_fail_safe(tmp_path: Path) -> None:
snap = build_snapshot(tmp_path / "does-not-exist.sqlite")
assert snap.ok is False
assert snap.error and "not found" in snap.error
# Still renders a friendly page rather than raising.
html_out = render_html(snap)
assert "No data" in html_out
def test_build_snapshot_never_writes(tmp_path: Path) -> None:
@ -268,43 +187,11 @@ def test_build_snapshot_never_writes(tmp_path: Path) -> None:
assert not missing.exists() # mode=ro did not create it
# --- pipeline-map + /api/state JSON contract tests ---------------------------
def test_render_html_includes_svg_map_and_poller() -> None:
"""The page must carry the inline SVG map and the inline JS poller."""
html_out = render_html(_sample_snapshot())
# Inline SVG map (hand-rolled, not fetched).
assert '<svg class="map"' in html_out
assert 'data-stage="clarify"' in html_out
# Inline poller that fetches the JSON sidecar — no external libraries.
assert "<script>" in html_out
assert "fetch('/api/state'" in html_out
assert "setInterval(poll" in html_out
# Tooltip mount + live clock the poller updates.
assert 'id="tip"' in html_out
assert 'id="clock"' in html_out
# Legend present.
assert "awaiting human" in html_out
# Still fully offline: no CDN / external references.
assert "http://" not in html_out and "https://" not in html_out
# The no-JS fallback is the meta-refresh, now inside <noscript>.
assert "<noscript>" in html_out
assert '<meta http-equiv="refresh" content="10">' in html_out
def test_render_html_renders_map_even_with_no_data() -> None:
"""A not-ok snapshot still renders the map + poller so it goes live later."""
snap = Snapshot(generated_at="now", db_path="/x.sqlite", ok=False, error="boom")
html_out = render_html(snap)
assert '<svg class="map"' in html_out
assert "fetch('/api/state'" in html_out
assert "No data" in html_out
assert html_out.rstrip().endswith("</html>")
# --- /api/state JSON payload (snapshot_to_dict) contract ---------------------
def test_snapshot_to_dict_shape_and_grouping() -> None:
"""The /api/state payload has the expected keys + correct stage grouping."""
"""The /api/state payload keeps its keys + correct stage grouping."""
payload = snapshot_to_dict(_sample_snapshot())
assert payload["ok"] is True
assert set(payload) >= {
@ -317,32 +204,22 @@ def test_snapshot_to_dict_shape_and_grouping() -> None:
"summary",
"budget",
}
# One stage entry per declared STAGE, in declared order.
assert [s["key"] for s in payload["stages"]] == [s.key for s in STAGES]
by_key = {s["key"]: s for s in payload["stages"]}
# The waiting clarify task lands on the clarify node and flags awaiting_human.
clarify = by_key["clarify"]
assert clarify["count"] == 1
assert clarify["state"] == "awaiting_human"
assert clarify["tasks"][0]["short_id"] == "abc123de"
# Each stage carries its model/agent role for the per-node label.
assert clarify["agent"] == "Claude (sub)"
assert by_key["review"]["agent"] == "GPT-4.1 (cross)"
assert by_key["build"]["agent"] == "DeepSeek (fast)"
# The active plan task lands on the plan node as 'active'.
assert by_key["plan"]["state"] == "active"
assert by_key["plan"]["count"] == 1
# The parked task maps onto the (terminal-ish) parked side — no stage owns
# the 'parked' phase, so no pipeline node claims it.
# The parked task maps onto no pipeline node.
assert all(
t["thread_id"] != "dead0000beef" for s in payload["stages"] for t in s["tasks"]
)
# An idle stage with no tasks.
assert by_key["intake"]["state"] == "idle"
assert by_key["intake"]["count"] == 0
# Summary mirrors the snapshot counts.
assert payload["summary"]["active"] == 2
assert payload["summary"]["waiting"] == 1
assert payload["summary"]["open_questions"] == 1
@ -350,20 +227,17 @@ def test_snapshot_to_dict_shape_and_grouping() -> None:
def test_snapshot_to_dict_not_ok_is_serializable() -> None:
"""A not-ok snapshot still serializes to a valid, JSON-dumpable payload."""
snap = Snapshot(generated_at="now", db_path="/x.sqlite", ok=False, error="boom")
payload = snapshot_to_dict(snap)
assert payload["ok"] is False
assert payload["error"] == "boom"
# Stages still present (all idle) so the client can render the empty map.
assert [s["key"] for s in payload["stages"]] == [s.key for s in STAGES]
assert all(s["count"] == 0 for s in payload["stages"])
# Round-trips through json.
assert json.loads(json.dumps(payload))["ok"] is False
def test_descriptions_escaped_in_html_and_json_seed() -> None:
"""Hostile descriptions must not break out of HTML *or* the inline JSON."""
def test_snapshot_to_dict_carries_raw_description_for_textcontent() -> None:
"""The JSON sidecar carries the raw description; the SPA renders it safely."""
snap = Snapshot(
generated_at="now",
db_path="/x.sqlite",
@ -378,39 +252,6 @@ def test_descriptions_escaped_in_html_and_json_seed() -> None:
)
],
)
html_out = render_html(snap)
# No raw script tag survives anywhere in the document (table render escapes
# it; the inline JSON seed escapes < and > to \\uXXXX).
assert "<script>alert(1)" not in html_out
assert "</script><script>" not in html_out
assert "\\u003cscript\\u003e" in html_out # JSON seed escaped form
# The JSON sidecar itself is valid JSON carrying the raw description (the
# consumer renders it via textContent, never as markup).
payload = snapshot_to_dict(snap)
body = json.dumps(payload)
parsed = json.loads(body)
parsed = json.loads(json.dumps(payload))
assert parsed["tasks"][0]["task"] == "</script><script>alert(1)</script>"
def test_build_snapshot_api_payload_against_seeded_ledger(tmp_path: Path) -> None:
"""End-to-end: read a seeded ledger, then assert the /api/state payload."""
db = tmp_path / "agent_team.sqlite"
_seed_db(db)
payload = snapshot_to_dict(build_snapshot(db))
assert payload["ok"] is True
by_key = {s["key"]: s for s in payload["stages"]}
# thread-waiting (clarify, open question) -> clarify node, awaiting_human.
assert by_key["clarify"]["state"] == "awaiting_human"
assert by_key["clarify"]["count"] == 1
assert by_key["clarify"]["tasks"][0]["thread_id"] == "thread-waiting"
# thread-active (plan, active) -> plan node, active.
assert by_key["plan"]["state"] == "active"
assert by_key["plan"]["count"] == 1
# Budget surfaced.
assert payload["budget"]["available"] is True
assert payload["budget"]["total_usd"] == pytest.approx(0.25)
# Whole payload is JSON-serializable.
assert json.loads(json.dumps(payload))["summary"]["total"] == 2

View file

@ -0,0 +1,60 @@
"""Tests for the LangGraph-introspected pipeline topology (``agent_team.topology``)."""
from __future__ import annotations
from agent_team.topology import (
NODE_META,
PHASE_TO_NODE,
build_topology,
missing_meta,
node_for_phase,
)
def test_every_graph_node_has_meta() -> None:
"""A newly added graph node must fail loudly until its NODE_META is filled in."""
assert missing_meta() == []
def test_topology_has_expected_nodes_and_trees() -> None:
topo = build_topology()
ids = {n["id"] for n in topo["nodes"]}
# The maximal P3+ wiring nodes, derived from the real graph.
assert {"intake", "clarify", "plan", "review", "build_node", "verify_node"} <= ids
# __start__/__end__ are not drawn.
assert "__start__" not in ids and "__end__" not in ids
# intake roots the core tree; the SDLC pipeline branches off it.
tree_ids = {t["id"] for t in topo["trees"]}
assert {"core", "sdlc"} <= tree_ids
by_id = {n["id"]: n for n in topo["nodes"]}
assert by_id["intake"]["tree"] == "core"
assert by_id["clarify"]["kind"] == "gate" # human gate lives at clarify
# P3 nodes render but are flagged inert.
assert by_id["build_node"]["gated"] is True
def test_edges_classified_spine_branch_loopback() -> None:
topo = build_topology()
kinds = {(e["from"], e["to"]): e["kind"] for e in topo["edges"]}
assert kinds[("intake", "clarify")] == "spine"
assert kinds[("clarify", "plan")] == "spine"
# review -> plan is the revise loop; verify -> build_node is the build loop.
assert kinds[("review", "plan")] == "loopback"
assert kinds[("verify_node", "build_node")] == "loopback"
# No edge references the framework terminals.
nodes = {n["id"] for n in topo["nodes"]}
for e in topo["edges"]:
assert e["from"] in nodes and e["to"] in nodes
def test_phase_to_node_mapping() -> None:
assert node_for_phase("build") == "build_node"
assert node_for_phase("verify") == "verify_node"
assert node_for_phase("clarify") == "clarify"
# Terminal/exception phases map to no node.
assert node_for_phase("done") is None
assert node_for_phase("parked") is None
assert node_for_phase("") is None
# Every mapped target is a real meta node.
for node_id in PHASE_TO_NODE.values():
assert node_id in NODE_META

View file

@ -0,0 +1,163 @@
"""Tests for the task_transitions recorder + graph instrumentation.
Covers (plan refs): idempotent record_entry under resume replay (BLOCK-1), the
terminal close that fills exited_at on the last node (N2), fail-soft writes that
never break the pipeline, the B4 signature-preservation guarantee of the
``_instrument`` wrapper, and an end-to-end graph drive that records rows.
"""
from __future__ import annotations
import inspect
import sqlite3
from pathlib import Path
from agent_team.db import connect, init_db, read_transitions
from agent_team.db.schema import assert_task_transitions_ready
from agent_team.db.transitions import TransitionRecorder
def _rows(db: Path, thread_id: str) -> list[dict[str, object]]:
conn = connect(db)
try:
return read_transitions(conn, thread_id)
finally:
conn.close()
def test_record_entry_chains_and_closes_previous(tmp_path: Path) -> None:
db = tmp_path / "led.sqlite"
init_db(db)
rec = TransitionRecorder(db)
rec.record_entry(thread_id="t1", to_phase="intake", status="active")
rec.record_entry(thread_id="t1", to_phase="clarify", status="active")
rows = _rows(db, "t1")
assert [r["to_phase"] for r in rows] == ["intake", "clarify"]
# First row closed (exited_at set) when the second landed; chain carried.
assert rows[0]["exited_at"] is not None
assert rows[1]["exited_at"] is None
assert rows[1]["from_phase"] == "intake"
def test_record_entry_is_idempotent_under_replay(tmp_path: Path) -> None:
"""Re-entering the SAME open node (resume replay) does not duplicate a row."""
db = tmp_path / "led.sqlite"
init_db(db)
rec = TransitionRecorder(db)
rec.record_entry(thread_id="t1", to_phase="clarify", status="active")
rec.record_entry(thread_id="t1", to_phase="clarify", status="active") # replay
rec.record_entry(thread_id="t1", to_phase="clarify", status="active") # replay
rows = _rows(db, "t1")
assert len(rows) == 1
assert rows[0]["to_phase"] == "clarify"
assert rows[0]["exited_at"] is None # still open, single row
def test_close_terminal_stamps_exit_on_last_open_row(tmp_path: Path) -> None:
db = tmp_path / "led.sqlite"
init_db(db)
rec = TransitionRecorder(db)
rec.record_entry(thread_id="t1", to_phase="plan", status="active")
rec.close_terminal(thread_id="t1", status="done")
rows = _rows(db, "t1")
assert len(rows) == 1
assert rows[0]["exited_at"] is not None
assert rows[0]["status"] == "done"
def test_writes_are_fail_soft(tmp_path: Path) -> None:
"""A ledger error is swallowed — recording must never raise into the graph."""
# Point the recorder at a path that is a directory, so connect/execute fails.
bad = tmp_path / "a_dir"
bad.mkdir()
rec = TransitionRecorder(bad)
# Must not raise.
rec.record_entry(thread_id="t1", to_phase="intake", status="active")
rec.close_terminal(thread_id="t1", status="done")
def test_read_transitions_missing_table_returns_empty(tmp_path: Path) -> None:
db = tmp_path / "empty.sqlite"
# A bare DB with no agent-team schema.
sqlite3.connect(str(db)).close()
conn = connect(db)
try:
assert read_transitions(conn, "t1") == []
finally:
conn.close()
def test_assert_task_transitions_ready(tmp_path: Path) -> None:
db = tmp_path / "led.sqlite"
init_db(db)
# Does not raise on a properly migrated DB.
assert_task_transitions_ready(db)
# --- B4: _instrument signature preservation + graph integration --------------
def test_instrument_preserves_signature_and_forwards_config() -> None:
"""The wrapper must present fn's signature (so LangGraph injects config) and
forward every argument verbatim (plan B4)."""
from agent_team.graph import _instrument
seen: dict[str, object] = {}
def node_with_config(state: dict, config: dict) -> dict: # type: ignore[type-arg]
seen["state"] = state
seen["config"] = config
return {"status": "active"}
class _Rec:
def record_entry(self, **_kw: object) -> None: ...
def close_terminal(self, **_kw: object) -> None: ...
wrapped = _instrument("plan", node_with_config, _Rec())
# Signature still advertises BOTH params, so LangGraph passes config.
assert list(inspect.signature(wrapped).parameters) == ["state", "config"]
# And the wrapper forwards them.
wrapped({"thread_id": "t"}, {"configurable": {}})
assert seen["config"] == {"configurable": {}}
def test_instrument_none_recorder_is_identity() -> None:
from agent_team.graph import _instrument
def node(state: dict) -> dict: # type: ignore[type-arg]
return {}
assert _instrument("intake", node, None) is node
def test_graph_drive_records_transitions(tmp_path: Path) -> None:
"""End-to-end: a recorder-instrumented graph records rows as a task runs."""
from langgraph.checkpoint.memory import MemorySaver
from agent_team.graph import build_graph, resume_task, start_task
db = tmp_path / "led.sqlite"
init_db(db)
rec = TransitionRecorder(db)
graph = build_graph(MemorySaver(), transition_recorder=rec)
tid, _ = start_task(graph, transport="cli", task="do a thing")
# Ran INTAKE -> CLARIFY (suspended on the human gate).
rows = _rows(db, tid)
phases = [r["to_phase"] for r in rows]
assert "intake" in phases
assert "clarify" in phases
resume_task(graph, thread_id=tid, answer="scope it")
rows = _rows(db, tid)
phases = [r["to_phase"] for r in rows]
assert "plan" in phases
# Terminal plan node (status DONE) closed its row.
plan_row = [r for r in rows if r["to_phase"] == "plan"][-1]
assert plan_row["exited_at"] is not None

7
agent-team/web/.gitignore vendored Normal file
View file

@ -0,0 +1,7 @@
node_modules/
dist/
*.tsbuildinfo
.vite/
coverage/
playwright-report/
test-results/

13
agent-team/web/index.html Normal file
View file

@ -0,0 +1,13 @@
<!doctype html>
<html lang="en">
<head>
<meta charset="UTF-8" />
<meta name="viewport" content="width=device-width, initial-scale=1.0" />
<meta name="color-scheme" content="dark" />
<title>agent-team · pipeline</title>
</head>
<body>
<div id="root"></div>
<script type="module" src="/src/main.tsx"></script>
</body>
</html>

4290
agent-team/web/package-lock.json generated Normal file

File diff suppressed because it is too large Load diff

View file

@ -0,0 +1,33 @@
{
"name": "agent-team-dashboard",
"private": true,
"version": "1.0.0",
"type": "module",
"description": "Read-only status dashboard SPA for the Sea Haven agent-team pipeline.",
"scripts": {
"dev": "vite",
"build": "tsc -b && vite build",
"preview": "vite preview",
"test": "vitest run",
"test:watch": "vitest",
"typecheck": "tsc -b --noEmit"
},
"dependencies": {
"@dagrejs/dagre": "1.1.4",
"react": "18.3.1",
"react-dom": "18.3.1",
"reactflow": "11.11.4"
},
"devDependencies": {
"@testing-library/jest-dom": "6.4.8",
"@testing-library/react": "16.0.1",
"@types/node": "20.14.15",
"@types/react": "18.3.5",
"@types/react-dom": "18.3.0",
"@vitejs/plugin-react": "4.3.1",
"jsdom": "24.1.3",
"typescript": "5.5.4",
"vite": "5.4.8",
"vitest": "1.6.0"
}
}

129
agent-team/web/src/App.tsx Normal file
View file

@ -0,0 +1,129 @@
import { useCallback, useEffect, useMemo, useRef, useState } from "react";
import {
fetchState,
fetchTask,
fetchTopology,
type StateResponse,
type TaskDetail,
type Topology,
} from "./api";
import { PipelineMap } from "./components/PipelineMap";
import { TaskDrawer } from "./components/TaskDrawer";
import { TaskList } from "./components/TaskList";
import { TopBar } from "./components/TopBar";
const POLL_MS = 4000;
export default function App() {
const [state, setState] = useState<StateResponse | null>(null);
const [topology, setTopology] = useState<Topology | null>(null);
const [connected, setConnected] = useState(false);
const [selectedId, setSelectedId] = useState<string | null>(null);
const [detail, setDetail] = useState<TaskDetail | null>(null);
const [detailLoading, setDetailLoading] = useState(false);
const [detailError, setDetailError] = useState<string | null>(null);
const [nodeFilter, setNodeFilter] = useState<string | null>(null);
// Topology is static for the life of the page — fetch once.
useEffect(() => {
fetchTopology()
.then(setTopology)
.catch(() => setTopology({ trees: [], nodes: [], edges: [] }));
}, []);
// Poll live state every 4s; the connection dot reflects the last fetch.
useEffect(() => {
let alive = true;
const tick = async () => {
try {
const s = await fetchState();
if (alive) {
setState(s);
setConnected(true);
}
} catch {
if (alive) setConnected(false);
}
};
tick();
const id = setInterval(tick, POLL_MS);
return () => {
alive = false;
clearInterval(id);
};
}, []);
const detailReq = useRef(0);
const selectTask = useCallback((threadId: string) => {
setSelectedId(threadId);
setDetailLoading(true);
setDetailError(null);
const req = ++detailReq.current;
fetchTask(threadId)
.then((d) => {
if (req === detailReq.current) setDetail(d);
})
.catch((e: unknown) => {
if (req === detailReq.current)
setDetailError(e instanceof Error ? e.message : "load failed");
})
.finally(() => {
if (req === detailReq.current) setDetailLoading(false);
});
}, []);
const closeDrawer = useCallback(() => {
setSelectedId(null);
setDetail(null);
setDetailError(null);
}, []);
const toggleNodeFilter = useCallback((nodeId: string) => {
setNodeFilter((cur) => (cur === nodeId ? null : nodeId));
}, []);
// Highlight the nodes a selected task has visited.
const pathNodeIds = useMemo(() => {
if (!selectedId || !detail?.ok) return null;
return new Set(detail.timeline.map((s) => s.to_phase));
}, [selectedId, detail]);
return (
<div className="app">
<TopBar
state={state}
connected={connected}
generatedAt={state?.generated_at ?? null}
/>
<div className="body">
<TaskList
state={state}
topology={topology}
selectedId={selectedId}
nodeFilter={nodeFilter}
onSelect={selectTask}
onClearNodeFilter={() => setNodeFilter(null)}
/>
<main className="center">
{topology ? (
<PipelineMap
topology={topology}
state={state}
pathNodeIds={pathNodeIds}
onNodeClick={toggleNodeFilter}
/>
) : (
<div className="loading">Loading pipeline…</div>
)}
</main>
<TaskDrawer
detail={detail}
loading={detailLoading}
error={detailError}
onClose={closeDrawer}
/>
</div>
</div>
);
}

102
agent-team/web/src/api.ts Normal file
View file

@ -0,0 +1,102 @@
// Typed client for the read-only dashboard API (agent_team/dashboard.py).
// All endpoints are GET + read-only; the SPA never mutates server state.
export type NodeState = "idle" | "active" | "awaiting_human" | "parked";
export interface TaskView {
thread_id: string;
short_id: string;
task: string;
current_phase: string;
status: string;
waiting: boolean;
waiting_since: string | null;
}
export interface StateResponse {
ok: boolean;
error: string | null;
generated_at: string;
db_path: string;
tasks: TaskView[];
waiting: TaskView[];
question_counts: Record<string, number>;
summary: {
active: number;
parked: number;
waiting: number;
open_questions: number;
total: number;
};
budget: {
available: boolean;
total_usd?: number | null;
recent?: Array<Record<string, unknown>>;
};
// Per-topology-node live state, keyed by node id. Absent nodes are idle.
nodes: Record<string, { state: NodeState; count: number }>;
}
export interface TopoNode {
id: string;
label: string;
agent: string;
tree: string;
kind: string;
gated: boolean;
}
export interface TopoEdge {
from: string;
to: string;
kind: "spine" | "branch" | "loopback";
}
export interface Topology {
trees: Array<{ id: string; label: string; root?: boolean }>;
nodes: TopoNode[];
edges: TopoEdge[];
}
export interface TimelineStep {
from_phase: string | null;
to_phase: string;
entered_at: string | null;
exited_at: string | null;
duration_s: number | null;
status: string | null;
note: string | null;
cost_usd: number;
}
export interface TaskDetail {
ok: boolean;
error?: string;
thread_id: string;
short_id: string;
task: string;
status: string;
current_phase: string;
partial: boolean;
timeline: TimelineStep[];
qa_history: unknown[];
review_verdicts: unknown[];
plan: Record<string, unknown> | null;
cost_by_stage: Record<string, number>;
total_usd: number;
}
async function getJSON<T>(url: string): Promise<T> {
const res = await fetch(url, { headers: { Accept: "application/json" } });
if (!res.ok) throw new Error(`${url} -> ${res.status}`);
return (await res.json()) as T;
}
export const fetchState = (): Promise<StateResponse> =>
getJSON<StateResponse>(`/api/state?ts=${Date.now()}`);
export const fetchTopology = (): Promise<Topology> =>
getJSON<Topology>("/api/topology");
export const fetchTask = (threadId: string): Promise<TaskDetail> =>
getJSON<TaskDetail>(`/api/task/${encodeURIComponent(threadId)}?ts=${Date.now()}`);

View file

@ -0,0 +1,48 @@
import { memo } from "react";
import { Handle, Position, type NodeProps } from "reactflow";
import type { RFData } from "../layout";
// One pipeline stage on the map. Color encodes live state; a badge shows the
// live task count; gated (P3/inert) nodes render dimmed with a tag. ARIA role +
// label keep the map reachable for screen readers / keyboard users.
function MapNodeImpl({ data }: NodeProps<RFData>) {
const { meta, state, count, dimmed, onPath } = data;
const cls = [
"node",
`state-${state}`,
meta.gated ? "gated" : "",
dimmed ? "dimmed" : "",
onPath ? "on-path" : "",
]
.filter(Boolean)
.join(" ");
return (
<div
className={cls}
role="group"
aria-label={`${meta.label} stage, ${state}, ${count} task${count === 1 ? "" : "s"}`}
tabIndex={0}
>
<Handle type="target" position={Position.Left} />
<div className="node-label">
{meta.label}
{meta.kind === "gate" && (
<span className="gate-pip" title="human gate" aria-hidden>
⛋
</span>
)}
</div>
<div className="node-agent">{meta.agent}</div>
{count > 0 && (
<span className="node-count" aria-hidden>
{count}
</span>
)}
{meta.gated && <span className="node-gated-tag">inert</span>}
<Handle type="source" position={Position.Right} />
</div>
);
}
export const MapNode = memo(MapNodeImpl);

View file

@ -0,0 +1,55 @@
import { useMemo } from "react";
import ReactFlow, {
Background,
Controls,
type NodeMouseHandler,
ReactFlowProvider,
} from "reactflow";
import "reactflow/dist/style.css";
import type { StateResponse, Topology } from "../api";
import { layout } from "../layout";
import { MapNode } from "./MapNode";
const nodeTypes = { pipeline: MapNode };
interface Props {
topology: Topology;
state: StateResponse | null;
pathNodeIds: Set<string> | null;
onNodeClick: (nodeId: string) => void;
}
// The center pipeline map. Nodes are dagre-laid from the topology; live state
// colors them; clicking a node filters the task list. Multiple trees (processes
// off intake) lay out as separate lanes automatically.
export function PipelineMap({ topology, state, pathNodeIds, onNodeClick }: Props) {
const { nodes, edges } = useMemo(
() => layout(topology, state?.nodes ?? {}, pathNodeIds),
[topology, state, pathNodeIds],
);
const handleNodeClick: NodeMouseHandler = (_evt, node) => onNodeClick(node.id);
return (
<div className="map-pane" aria-label="pipeline map">
<ReactFlowProvider>
<ReactFlow
nodes={nodes}
edges={edges}
nodeTypes={nodeTypes}
fitView
nodesDraggable={false}
nodesConnectable={false}
elementsSelectable
onNodeClick={handleNodeClick}
proOptions={{ hideAttribution: true }}
minZoom={0.3}
maxZoom={1.6}
>
<Background color="#222a36" gap={22} />
<Controls showInteractive={false} />
</ReactFlow>
</ReactFlowProvider>
</div>
);
}

View file

@ -0,0 +1,53 @@
import { render, screen } from "@testing-library/react";
import { describe, expect, it } from "vitest";
import { TaskDrawer } from "./TaskDrawer";
import type { TaskDetail } from "../api";
const DETAIL: TaskDetail = {
ok: true,
thread_id: "aaaa1111bbbb",
short_id: "aaaa1111",
task: "remediate CVE-2026-1",
status: "done",
current_phase: "done",
partial: false,
timeline: [
{ from_phase: null, to_phase: "intake", entered_at: "2026-06-23T00:00:00+00:00", exited_at: "2026-06-23T00:00:01+00:00", duration_s: 1, status: "active", note: null, cost_usd: 0 },
{ from_phase: "intake", to_phase: "plan", entered_at: "2026-06-23T00:00:01+00:00", exited_at: null, duration_s: null, status: "active", note: null, cost_usd: 0.12 },
],
qa_history: [{ turn: 0, answer: "scope it" }],
review_verdicts: [],
plan: { summary: "do the thing" },
cost_by_stage: { plan: 0.12 },
total_usd: 0.12,
};
describe("TaskDrawer", () => {
it("renders the timeline steps with the task description and spend", () => {
render(<TaskDrawer detail={DETAIL} loading={false} error={null} onClose={() => {}} />);
expect(screen.getByText("remediate CVE-2026-1")).toBeInTheDocument();
expect(screen.getByText("intake")).toBeInTheDocument();
expect(screen.getByText("plan")).toBeInTheDocument();
expect(screen.getByText("in progress")).toBeInTheDocument(); // open final step
expect(screen.getByText(/spend \$0\.12/)).toBeInTheDocument();
});
it("shows a partial-history banner when partial", () => {
render(
<TaskDrawer
detail={{ ...DETAIL, partial: true }}
loading={false}
error={null}
onClose={() => {}}
/>,
);
expect(screen.getByText(/Partial history/)).toBeInTheDocument();
});
it("renders nothing when idle (no detail, not loading)", () => {
const { container } = render(
<TaskDrawer detail={null} loading={false} error={null} onClose={() => {}} />,
);
expect(container).toBeEmptyDOMElement();
});
});

View file

@ -0,0 +1,109 @@
import type { TaskDetail } from "../api";
interface Props {
detail: TaskDetail | null;
loading: boolean;
error: string | null;
onClose: () => void;
}
function fmtDuration(s: number | null): string {
if (s == null) return "—";
if (s < 60) return `${s.toFixed(1)}s`;
if (s < 3600) return `${(s / 60).toFixed(1)}m`;
return `${(s / 3600).toFixed(1)}h`;
}
function fmtCost(usd: number): string {
return usd > 0 ? `$${usd.toFixed(4)}` : "—";
}
// Right-hand drawer: a task's journey through the pipeline (vertical timeline of
// task_transitions) plus Q&A, review verdicts, plan summary, and total spend.
export function TaskDrawer({ detail, loading, error, onClose }: Props) {
if (!detail && !loading && !error) return null;
return (
<section className="drawer" aria-label="task detail" role="complementary">
<div className="drawer-head">
<h2>Task history</h2>
<button className="drawer-close" aria-label="Close detail" onClick={onClose}>
✕
</button>
</div>
{loading && <div className="drawer-loading">Loading…</div>}
{error && <div className="drawer-error">{error}</div>}
{detail && detail.ok && (
<>
<div className="drawer-summary">
<div className="drawer-task">{detail.task}</div>
<div className="drawer-meta">
<span className={`badge ${detail.status}`}>{detail.status}</span>
<span className="mono">{detail.short_id}</span>
<span className="spend">spend {fmtCost(detail.total_usd)}</span>
</div>
{detail.partial && (
<div className="partial-banner">
Partial history — recorded since the v3 ledger; earlier steps are
reconstructed.
</div>
)}
</div>
<ol className="timeline">
{detail.timeline.map((step, i) => {
const open = step.exited_at == null;
return (
<li className={`tl-step ${open ? "open" : ""}`} key={i}>
<span className="tl-dot" aria-hidden />
<div className="tl-body">
<div className="tl-node">{step.to_phase}</div>
<div className="tl-times">
<span>{step.entered_at ?? "—"}</span>
<span className="tl-dur">
{open ? "in progress" : fmtDuration(step.duration_s)}
</span>
</div>
<div className="tl-foot">
{step.status && <span className={`badge ${step.status}`}>{step.status}</span>}
<span className="tl-cost">{fmtCost(step.cost_usd)}</span>
{step.note && <span className="tl-note">{step.note}</span>}
</div>
</div>
</li>
);
})}
{detail.timeline.length === 0 && (
<li className="empty">No transitions recorded.</li>
)}
</ol>
{detail.qa_history.length > 0 && (
<details className="drawer-section" open>
<summary>Q&amp;A ({detail.qa_history.length})</summary>
<pre className="json">{JSON.stringify(detail.qa_history, null, 2)}</pre>
</details>
)}
{detail.review_verdicts.length > 0 && (
<details className="drawer-section">
<summary>Review verdicts ({detail.review_verdicts.length})</summary>
<pre className="json">{JSON.stringify(detail.review_verdicts, null, 2)}</pre>
</details>
)}
{detail.plan && (
<details className="drawer-section">
<summary>Plan</summary>
<pre className="json">{JSON.stringify(detail.plan, null, 2)}</pre>
</details>
)}
</>
)}
{detail && !detail.ok && (
<div className="drawer-error">{detail.error ?? "Could not load task."}</div>
)}
</section>
);
}

View file

@ -0,0 +1,74 @@
import { fireEvent, render, screen } from "@testing-library/react";
import { describe, expect, it, vi } from "vitest";
import { TaskList } from "./TaskList";
import type { StateResponse, Topology } from "../api";
const TOPO: Topology = {
trees: [
{ id: "core", label: "Core" },
{ id: "sdlc", label: "SDLC" },
],
nodes: [
{ id: "intake", label: "Intake", agent: "coordinator", tree: "core", kind: "phase", gated: false },
{ id: "plan", label: "Plan", agent: "Claude", tree: "sdlc", kind: "phase", gated: false },
],
edges: [],
};
function makeState(): StateResponse {
return {
ok: true,
error: null,
generated_at: "now",
db_path: "x",
tasks: [
{ thread_id: "aaaa1111", short_id: "aaaa1111", task: "remediate CVE", current_phase: "plan", status: "active", waiting: false, waiting_since: null },
{ thread_id: "bbbb2222", short_id: "bbbb2222", task: "fix login bug", current_phase: "intake", status: "parked", waiting: false, waiting_since: null },
],
waiting: [],
question_counts: {},
summary: { active: 1, parked: 1, waiting: 0, open_questions: 0, total: 2 },
budget: { available: false },
nodes: {},
};
}
describe("TaskList", () => {
it("renders all tasks", () => {
render(
<TaskList state={makeState()} topology={TOPO} selectedId={null} nodeFilter={null}
onSelect={() => {}} onClearNodeFilter={() => {}} />,
);
expect(screen.getByText("remediate CVE")).toBeInTheDocument();
expect(screen.getByText("fix login bug")).toBeInTheDocument();
});
it("filters by status", () => {
render(
<TaskList state={makeState()} topology={TOPO} selectedId={null} nodeFilter={null}
onSelect={() => {}} onClearNodeFilter={() => {}} />,
);
fireEvent.change(screen.getByLabelText("Filter by status"), { target: { value: "parked" } });
expect(screen.queryByText("remediate CVE")).not.toBeInTheDocument();
expect(screen.getByText("fix login bug")).toBeInTheDocument();
});
it("honors a node filter (map click)", () => {
render(
<TaskList state={makeState()} topology={TOPO} selectedId={null} nodeFilter="plan"
onSelect={() => {}} onClearNodeFilter={() => {}} />,
);
expect(screen.getByText("remediate CVE")).toBeInTheDocument();
expect(screen.queryByText("fix login bug")).not.toBeInTheDocument();
});
it("fires onSelect when a row is clicked", () => {
const onSelect = vi.fn();
render(
<TaskList state={makeState()} topology={TOPO} selectedId={null} nodeFilter={null}
onSelect={onSelect} onClearNodeFilter={() => {}} />,
);
fireEvent.click(screen.getByText("remediate CVE"));
expect(onSelect).toHaveBeenCalledWith("aaaa1111");
});
});

View file

@ -0,0 +1,145 @@
import { useMemo, useState } from "react";
import type { StateResponse, TaskView, Topology } from "../api";
// Mirror of the backend phase->node mapping (topology.PHASE_TO_NODE) so the list
// can resolve a task's tree and honor a map-node filter. Identity except the P3
// build/verify vertices.
const PHASE_TO_NODE: Record<string, string> = {
build: "build_node",
verify: "verify_node",
};
const phaseToNode = (phase: string): string => PHASE_TO_NODE[phase] ?? phase;
interface Props {
state: StateResponse | null;
topology: Topology | null;
selectedId: string | null;
nodeFilter: string | null;
onSelect: (threadId: string) => void;
onClearNodeFilter: () => void;
}
const STATUS_OPTIONS = ["all", "active", "waiting_human", "parked", "done", "failed"];
export function TaskList({
state,
topology,
selectedId,
nodeFilter,
onSelect,
onClearNodeFilter,
}: Props) {
const [query, setQuery] = useState("");
const [status, setStatus] = useState("all");
const [tree, setTree] = useState("all");
const treeByNode = useMemo(() => {
const m: Record<string, string> = {};
for (const n of topology?.nodes ?? []) m[n.id] = n.tree;
return m;
}, [topology]);
const tasks = state?.tasks ?? [];
const filtered = useMemo(() => {
const q = query.trim().toLowerCase();
return tasks.filter((t) => {
if (nodeFilter && phaseToNode(t.current_phase) !== nodeFilter) return false;
if (status !== "all" && t.status !== status) return false;
if (tree !== "all" && treeByNode[phaseToNode(t.current_phase)] !== tree)
return false;
if (q && !(`${t.task} ${t.short_id}`.toLowerCase().includes(q))) return false;
return true;
});
}, [tasks, query, status, tree, nodeFilter, treeByNode]);
return (
<aside className="tasklist" aria-label="task list">
<div className="tasklist-controls">
<input
className="search"
type="search"
placeholder="Search tasks…"
aria-label="Search tasks"
value={query}
onChange={(e) => setQuery(e.target.value)}
/>
<div className="filters">
<select
aria-label="Filter by status"
value={status}
onChange={(e) => setStatus(e.target.value)}
>
{STATUS_OPTIONS.map((s) => (
<option key={s} value={s}>
{s === "all" ? "All statuses" : s.replace("_", " ")}
</option>
))}
</select>
<select
aria-label="Filter by tree"
value={tree}
onChange={(e) => setTree(e.target.value)}
>
<option value="all">All trees</option>
{(topology?.trees ?? []).map((t) => (
<option key={t.id} value={t.id}>
{t.label}
</option>
))}
</select>
</div>
{nodeFilter && (
<button className="filter-chip" onClick={onClearNodeFilter}>
node: {nodeFilter} ✕
</button>
)}
</div>
<ul className="tasks" role="listbox" aria-label="tasks">
{filtered.length === 0 && <li className="empty">No tasks match.</li>}
{filtered.map((t) => (
<TaskRow
key={t.thread_id}
task={t}
selected={t.thread_id === selectedId}
onSelect={onSelect}
/>
))}
</ul>
</aside>
);
}
function TaskRow({
task,
selected,
onSelect,
}: {
task: TaskView;
selected: boolean;
onSelect: (id: string) => void;
}) {
return (
<li
className={`task-row ${selected ? "selected" : ""}`}
role="option"
aria-selected={selected}
tabIndex={0}
onClick={() => onSelect(task.thread_id)}
onKeyDown={(e) => {
if (e.key === "Enter" || e.key === " ") {
e.preventDefault();
onSelect(task.thread_id);
}
}}
>
<div className="task-row-top">
<span className={`badge ${task.status}`}>{task.status.replace("_", " ")}</span>
<span className="task-phase">{task.current_phase}</span>
{task.waiting && <span className="task-wait" title="awaiting human">⏳</span>}
</div>
<div className="task-desc">{task.task}</div>
<div className="task-id mono">{task.short_id}</div>
</li>
);
}

View file

@ -0,0 +1,42 @@
import type { StateResponse } from "../api";
interface Props {
state: StateResponse | null;
connected: boolean;
generatedAt: string | null;
}
const CHIPS: Array<{ key: keyof StateResponse["summary"]; label: string; cls: string }> =
[
{ key: "active", label: "Active", cls: "active" },
{ key: "waiting", label: "Awaiting human", cls: "waiting" },
{ key: "parked", label: "Parked", cls: "parked" },
{ key: "open_questions", label: "Open questions", cls: "open" },
{ key: "total", label: "Total", cls: "total" },
];
// Top status bar: summary chips, a connection/health dot, and the snapshot clock.
export function TopBar({ state, connected, generatedAt }: Props) {
const summary = state?.summary;
return (
<header className="topbar" role="banner">
<div className="brand">
<span className="brand-dot" aria-hidden />
<span className="brand-name">agent-team</span>
<span className="brand-sub">pipeline</span>
</div>
<div className="chips" role="list">
{CHIPS.map((c) => (
<div className={`chip chip-${c.cls}`} role="listitem" key={c.key}>
<span className="chip-num">{summary ? summary[c.key] : "–"}</span>
<span className="chip-lbl">{c.label}</span>
</div>
))}
</div>
<div className="health" title={connected ? "live" : "disconnected"}>
<span className={`health-dot ${connected ? "ok" : "down"}`} aria-hidden />
<span className="health-clock">{generatedAt ?? "—"}</span>
</div>
</header>
);
}

View file

@ -0,0 +1,59 @@
import { describe, expect, it } from "vitest";
import { layout } from "./layout";
import type { Topology } from "./api";
const TOPO: Topology = {
trees: [
{ id: "core", label: "Core", root: true },
{ id: "sdlc", label: "SDLC", root: false },
],
nodes: [
{ id: "intake", label: "Intake", agent: "coordinator", tree: "core", kind: "phase", gated: false },
{ id: "clarify", label: "Clarify", agent: "Claude", tree: "sdlc", kind: "gate", gated: false },
{ id: "plan", label: "Plan", agent: "Claude", tree: "sdlc", kind: "phase", gated: false },
{ id: "review", label: "Review", agent: "GPT-4.1", tree: "sdlc", kind: "phase", gated: false },
],
edges: [
{ from: "intake", to: "clarify", kind: "spine" },
{ from: "clarify", to: "plan", kind: "spine" },
{ from: "plan", to: "review", kind: "spine" },
{ from: "review", to: "plan", kind: "loopback" },
],
};
describe("layout", () => {
it("produces a laid-out node per topology node", () => {
const { nodes, edges } = layout(TOPO, {}, null);
expect(nodes).toHaveLength(4);
expect(edges).toHaveLength(4);
// every node has a numeric position assigned by dagre
for (const n of nodes) {
expect(typeof n.position.x).toBe("number");
expect(typeof n.position.y).toBe("number");
}
});
it("applies live state + count to the matching node", () => {
const { nodes } = layout(TOPO, { clarify: { state: "awaiting_human", count: 2 } }, null);
const clarify = nodes.find((n) => n.id === "clarify")!;
expect(clarify.data.state).toBe("awaiting_human");
expect(clarify.data.count).toBe(2);
const plan = nodes.find((n) => n.id === "plan")!;
expect(plan.data.state).toBe("idle");
});
it("marks the selected task's path and dims the rest", () => {
const path = new Set(["intake", "clarify"]);
const { nodes } = layout(TOPO, {}, path);
expect(nodes.find((n) => n.id === "intake")!.data.onPath).toBe(true);
expect(nodes.find((n) => n.id === "intake")!.data.dimmed).toBe(false);
expect(nodes.find((n) => n.id === "review")!.data.onPath).toBe(false);
expect(nodes.find((n) => n.id === "review")!.data.dimmed).toBe(true);
});
it("styles loopback edges distinctly (dashed)", () => {
const { edges } = layout(TOPO, {}, null);
const loop = edges.find((e) => e.id === "review->plan")!;
expect(loop.style?.strokeDasharray).toBeTruthy();
});
});

View file

@ -0,0 +1,71 @@
// Auto-layout the pipeline graph with dagre, so adding a node in the backend
// graph needs no manual coordinates here — the map redraws itself.
import dagre from "@dagrejs/dagre";
import type { Edge, Node } from "reactflow";
import type { Topology, TopoNode } from "./api";
export const NODE_W = 184;
export const NODE_H = 72;
export interface RFData {
meta: TopoNode;
state: "idle" | "active" | "awaiting_human" | "parked";
count: number;
dimmed: boolean;
onPath: boolean;
}
// Compute laid-out React Flow nodes + edges from the topology and live state.
// `onPath` marks the nodes a selected task has visited (highlighted); `dimmed`
// fades nodes outside the active task's path when one is selected.
export function layout(
topo: Topology,
liveNodes: Record<string, { state: RFData["state"]; count: number }>,
pathNodeIds: Set<string> | null,
): { nodes: Node<RFData>[]; edges: Edge[] } {
const g = new dagre.graphlib.Graph();
g.setGraph({ rankdir: "LR", nodesep: 36, ranksep: 72, marginx: 24, marginy: 24 });
g.setDefaultEdgeLabel(() => ({}));
for (const n of topo.nodes) g.setNode(n.id, { width: NODE_W, height: NODE_H });
for (const e of topo.edges) g.setEdge(e.from, e.to);
dagre.layout(g);
const nodes: Node<RFData>[] = topo.nodes.map((meta) => {
const pos = g.node(meta.id);
const live = liveNodes[meta.id];
const onPath = pathNodeIds ? pathNodeIds.has(meta.id) : false;
return {
id: meta.id,
type: "pipeline",
position: { x: pos.x - NODE_W / 2, y: pos.y - NODE_H / 2 },
data: {
meta,
state: live?.state ?? "idle",
count: live?.count ?? 0,
dimmed: pathNodeIds != null && !onPath,
onPath,
},
};
});
const edges: Edge[] = topo.edges.map((e) => {
const loop = e.kind === "loopback";
return {
id: `${e.from}->${e.to}`,
source: e.from,
target: e.to,
animated: false,
style: {
stroke: loop ? "#7a5cff" : e.kind === "branch" ? "#3f7fd8" : "#39414f",
strokeWidth: 1.5,
strokeDasharray: loop ? "5 4" : undefined,
},
label: loop ? "↺" : undefined,
labelStyle: { fill: "#9b8bff", fontWeight: 700 },
type: "smoothstep",
};
});
return { nodes, edges };
}

View file

@ -0,0 +1,10 @@
import { StrictMode } from "react";
import { createRoot } from "react-dom/client";
import App from "./App";
import "./theme.css";
createRoot(document.getElementById("root")!).render(
<StrictMode>
<App />
</StrictMode>,
);

View file

@ -0,0 +1,19 @@
import "@testing-library/jest-dom/vitest";
// React Flow measures the DOM; jsdom lacks ResizeObserver + layout APIs, so
// stub the minimum the component touches during tests.
class ResizeObserverStub {
observe() {}
unobserve() {}
disconnect() {}
}
globalThis.ResizeObserver =
globalThis.ResizeObserver ?? (ResizeObserverStub as typeof ResizeObserver);
if (!globalThis.DOMMatrixReadOnly) {
// @ts-expect-error - minimal stub for reactflow transforms in jsdom
globalThis.DOMMatrixReadOnly = class {
m22 = 1;
constructor() {}
};
}

View file

@ -0,0 +1,178 @@
:root {
--bg: #0c0e13;
--bg-1: #12151c;
--bg-2: #1a1e26;
--border: #232a36;
--border-2: #2f3645;
--text: #e6e8ee;
--text-dim: #8a93a6;
--accent: #7a5cff;
--accent-2: #3f7fd8;
--ok: #2f8f5e;
--ok-bg: #163a2b;
--ok-fg: #6ee7a8;
--warn-bg: #3a3216;
--warn-fg: #f5d76e;
--park-bg: #3a1f16;
--park-fg: #f3a36e;
--done-bg: #1c2b3a;
--done-fg: #7fb3f5;
--fail-bg: #3a1620;
--fail-fg: #f57f9c;
color-scheme: dark;
}
* { box-sizing: border-box; }
html, body, #root { height: 100%; margin: 0; }
body {
font-family: -apple-system, BlinkMacSystemFont, "Segoe UI", Roboto, Helvetica,
Arial, sans-serif;
background: var(--bg);
color: var(--text);
line-height: 1.45;
}
.mono { font-family: ui-monospace, SFMono-Regular, Menlo, monospace; }
.app { display: flex; flex-direction: column; height: 100%; }
.body { display: grid; grid-template-columns: 320px 1fr 360px; flex: 1; min-height: 0; }
/* --- Top bar -------------------------------------------------------------- */
.topbar {
display: flex; align-items: center; gap: 1.5rem;
padding: .6rem 1rem; background: var(--bg-1);
border-bottom: 1px solid var(--border);
}
.brand { display: flex; align-items: baseline; gap: .5rem; }
.brand-dot {
width: 10px; height: 10px; border-radius: 50%;
background: linear-gradient(135deg, var(--accent), var(--accent-2));
align-self: center;
}
.brand-name { font-weight: 700; letter-spacing: .01em; }
.brand-sub { color: var(--text-dim); font-size: .85rem; }
.chips { display: flex; gap: .6rem; flex: 1; }
.chip {
display: flex; flex-direction: column; align-items: center;
background: var(--bg-2); border: 1px solid var(--border);
border-radius: 8px; padding: .3rem .7rem; min-width: 4.5rem;
}
.chip-num { font-size: 1.25rem; font-weight: 700; }
.chip-lbl { font-size: .66rem; text-transform: uppercase; letter-spacing: .05em;
color: var(--text-dim); }
.chip-active .chip-num { color: var(--ok-fg); }
.chip-waiting .chip-num { color: var(--warn-fg); }
.chip-parked .chip-num { color: var(--park-fg); }
.health { display: flex; align-items: center; gap: .5rem; color: var(--text-dim);
font-size: .8rem; }
.health-dot { width: 8px; height: 8px; border-radius: 50%; }
.health-dot.ok { background: var(--ok-fg); box-shadow: 0 0 6px var(--ok); }
.health-dot.down { background: var(--fail-fg); }
/* --- Task list ------------------------------------------------------------ */
.tasklist { background: var(--bg-1); border-right: 1px solid var(--border);
display: flex; flex-direction: column; min-height: 0; }
.tasklist-controls { padding: .75rem; border-bottom: 1px solid var(--border);
display: flex; flex-direction: column; gap: .5rem; }
.search, select {
background: var(--bg-2); color: var(--text); border: 1px solid var(--border-2);
border-radius: 6px; padding: .4rem .55rem; font-size: .85rem; width: 100%;
}
.filters { display: flex; gap: .5rem; }
.filter-chip { background: var(--bg-2); color: var(--accent); cursor: pointer;
border: 1px solid var(--accent); border-radius: 999px; padding: .2rem .6rem;
font-size: .75rem; align-self: flex-start; }
.tasks { list-style: none; margin: 0; padding: .5rem; overflow-y: auto; flex: 1; }
.task-row { background: var(--bg-2); border: 1px solid var(--border);
border-radius: 8px; padding: .55rem .65rem; margin-bottom: .5rem; cursor: pointer;
transition: border-color .15s, transform .05s; }
.task-row:hover { border-color: var(--border-2); }
.task-row:focus-visible { outline: 2px solid var(--accent); outline-offset: 1px; }
.task-row.selected { border-color: var(--accent); box-shadow: 0 0 0 1px var(--accent); }
.task-row-top { display: flex; align-items: center; gap: .5rem; margin-bottom: .25rem; }
.task-phase { color: var(--text-dim); font-size: .75rem; }
.task-wait { margin-left: auto; }
.task-desc { font-size: .85rem; overflow: hidden; text-overflow: ellipsis;
display: -webkit-box; -webkit-line-clamp: 2; -webkit-box-orient: vertical; }
.task-id { color: var(--text-dim); font-size: .72rem; margin-top: .25rem; }
.empty { color: var(--text-dim); font-style: italic; padding: .5rem; }
/* --- Badges -------------------------------------------------------------- */
.badge { display: inline-block; padding: .08rem .45rem; border-radius: 999px;
font-size: .72rem; font-weight: 600; }
.badge.active { background: var(--ok-bg); color: var(--ok-fg); }
.badge.waiting_human, .badge.waiting { background: var(--warn-bg); color: var(--warn-fg); }
.badge.parked { background: var(--park-bg); color: var(--park-fg); }
.badge.done { background: var(--done-bg); color: var(--done-fg); }
.badge.failed { background: var(--fail-bg); color: var(--fail-fg); }
.badge.unknown { background: var(--bg-2); color: var(--text-dim); }
/* --- Map ----------------------------------------------------------------- */
.center { min-height: 0; min-width: 0; position: relative; }
.map-pane { width: 100%; height: 100%; }
.loading { padding: 2rem; color: var(--text-dim); }
.react-flow__attribution { display: none; }
.node {
width: 184px; height: 72px; background: var(--bg-2);
border: 1.5px solid var(--border-2); border-radius: 10px;
padding: .5rem .6rem; display: flex; flex-direction: column; justify-content: center;
position: relative; transition: border-color .25s, background .25s, opacity .2s;
}
.node:focus-visible { outline: 2px solid var(--accent); outline-offset: 2px; }
.node-label { font-weight: 700; font-size: .92rem; display: flex; gap: .3rem;
align-items: center; }
.gate-pip { color: var(--warn-fg); font-size: .8rem; }
.node-agent { color: var(--text-dim); font-size: .72rem; margin-top: .15rem; }
.node-count { position: absolute; top: -8px; right: -8px; min-width: 20px; height: 20px;
border-radius: 999px; background: var(--accent); color: #0c0e13; font-weight: 800;
font-size: .72rem; display: flex; align-items: center; justify-content: center;
padding: 0 5px; }
.node-gated-tag { position: absolute; bottom: 4px; right: 8px; font-size: .58rem;
letter-spacing: .08em; text-transform: uppercase; color: var(--text-dim); }
.node.state-active { border-color: var(--ok); background: var(--ok-bg); }
.node.state-awaiting_human { border-color: var(--warn-fg); background: var(--warn-bg); }
.node.state-parked { border-color: var(--park-fg); background: var(--park-bg); }
.node.gated { opacity: .62; }
.node.dimmed { opacity: .28; }
.node.on-path { box-shadow: 0 0 0 2px var(--accent), 0 0 14px rgba(122, 92, 255, .4); }
/* --- Drawer -------------------------------------------------------------- */
.drawer { background: var(--bg-1); border-left: 1px solid var(--border);
display: flex; flex-direction: column; min-height: 0; overflow-y: auto; }
.drawer-head { display: flex; align-items: center; justify-content: space-between;
padding: .75rem 1rem; border-bottom: 1px solid var(--border); position: sticky;
top: 0; background: var(--bg-1); }
.drawer-head h2 { font-size: 1rem; margin: 0; color: #9fb4d6; }
.drawer-close { background: none; border: none; color: var(--text-dim);
font-size: 1rem; cursor: pointer; }
.drawer-close:focus-visible { outline: 2px solid var(--accent); }
.drawer-loading, .drawer-error, .drawer-section { padding: 1rem; }
.drawer-error { color: var(--fail-fg); }
.drawer-summary { padding: 1rem; border-bottom: 1px solid var(--border); }
.drawer-task { font-size: .95rem; margin-bottom: .5rem; }
.drawer-meta { display: flex; gap: .6rem; align-items: center; flex-wrap: wrap;
font-size: .8rem; color: var(--text-dim); }
.spend { margin-left: auto; color: var(--text); }
.partial-banner { margin-top: .6rem; background: var(--warn-bg); color: var(--warn-fg);
border-radius: 6px; padding: .5rem .6rem; font-size: .76rem; }
.timeline { list-style: none; margin: 0; padding: 1rem 1rem 0; }
.tl-step { display: flex; gap: .7rem; padding-bottom: 1rem; position: relative; }
.tl-step:not(:last-child)::before { content: ""; position: absolute; left: 5px;
top: 14px; bottom: 0; width: 2px; background: var(--border-2); }
.tl-dot { width: 12px; height: 12px; border-radius: 50%; background: var(--accent);
margin-top: 2px; flex-shrink: 0; z-index: 1; }
.tl-step.open .tl-dot { background: var(--warn-fg); box-shadow: 0 0 6px var(--warn-fg); }
.tl-body { flex: 1; }
.tl-node { font-weight: 700; font-size: .88rem; }
.tl-times { display: flex; justify-content: space-between; color: var(--text-dim);
font-size: .72rem; margin: .15rem 0; }
.tl-dur { color: var(--text); }
.tl-foot { display: flex; gap: .5rem; align-items: center; font-size: .72rem; }
.tl-cost { color: var(--ok-fg); }
.tl-note { color: var(--text-dim); font-style: italic; }
.drawer-section summary { cursor: pointer; color: #9fb4d6; font-weight: 600;
font-size: .85rem; }
.json { background: var(--bg); border: 1px solid var(--border); border-radius: 6px;
padding: .6rem; font-size: .72rem; overflow-x: auto; max-height: 16rem;
font-family: ui-monospace, SFMono-Regular, Menlo, monospace; }

View file

@ -0,0 +1,21 @@
{
"compilerOptions": {
"target": "ES2020",
"useDefineForClassFields": true,
"lib": ["ES2020", "DOM", "DOM.Iterable"],
"module": "ESNext",
"skipLibCheck": true,
"moduleResolution": "bundler",
"allowImportingTsExtensions": true,
"resolveJsonModule": true,
"isolatedModules": true,
"noEmit": true,
"jsx": "react-jsx",
"strict": true,
"noUnusedLocals": true,
"noUnusedParameters": true,
"noFallthroughCasesInSwitch": true,
"types": ["node", "vitest/globals", "@testing-library/jest-dom"]
},
"include": ["src", "vite.config.ts"]
}

View file

@ -0,0 +1,30 @@
/// <reference types="vitest/config" />
import { defineConfig } from "vite";
import react from "@vitejs/plugin-react";
// The SPA is served by the FastAPI dashboard (agent_team/dashboard.py) from
// web/dist at LAN :8770. In dev, proxy /api to a locally running dashboard so
// `npm run dev` talks to real data without CORS. No CDN — everything bundles
// locally (the dashboard is offline/LAN-only).
export default defineConfig({
plugins: [react()],
build: {
outDir: "dist",
emptyOutDir: true,
sourcemap: false,
},
server: {
proxy: {
"/api": {
target: process.env.AGENT_TEAM_DASH || "http://127.0.0.1:8770",
changeOrigin: true,
},
},
},
test: {
globals: true,
environment: "jsdom",
setupFiles: ["./src/test/setup.ts"],
css: false,
},
});