diff --git a/agent-team/README.md b/agent-team/README.md index 8219fe1..59fac93 100644 --- a/agent-team/README.md +++ b/agent-team/README.md @@ -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 diff --git a/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index 53e67ca..474e112 100644 --- a/agent-team/agent_team/coordinator.py +++ b/agent-team/agent_team/coordinator.py @@ -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, diff --git a/agent-team/agent_team/dashboard.py b/agent-team/agent_team/dashboard.py new file mode 100644 index 0000000..69f5711 --- /dev/null +++ b/agent-team/agent_team/dashboard.py @@ -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"}) + + @app.get("/api/task/{thread_id}") + def api_task(thread_id: str) -> JSONResponse: + if not _THREAD_ID_RE.match(thread_id): + raise HTTPException(status_code=400, detail="invalid thread_id") + payload = task_detail(resolved_db, thread_id) + return JSONResponse(payload, headers={"Cache-Control": "no-store"}) + + # Static SPA, mounted LAST so it never shadows the /api routes. Skipped when + # no build is present (dev/test) so the API is still reachable. + dist = _resolve_dist_dir() + if dist.is_dir(): + from fastapi.staticfiles import StaticFiles + + app.mount("/", StaticFiles(directory=str(dist), html=True), name="spa") + + return app + + +def serve( + *, + host: str = _DEFAULT_HOST, + port: int = _DEFAULT_PORT, + db_path: str | Path | None = None, +) -> None: # pragma: no cover - process entry point + """Start the uvicorn server for the read-only dashboard (LAN, :8770). + + This is the systemd ExecStart target replacing ``status_page.serve``. + """ + try: + import uvicorn + except ImportError as exc: + raise RuntimeError( + "uvicorn is unavailable; install it: pip install uvicorn" + ) from exc + uvicorn.run(make_dashboard_app(db_path=db_path), host=host, port=port) diff --git a/agent-team/agent_team/db/__init__.py b/agent-team/agent_team/db/__init__.py index c62f898..6491a94 100644 --- a/agent-team/agent_team/db/__init__.py +++ b/agent-team/agent_team/db/__init__.py @@ -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", ] diff --git a/agent-team/agent_team/db/schema.py b/agent-team/agent_team/db/schema.py index 8b7cf47..35f1dda 100644 --- a/agent-team/agent_team/db/schema.py +++ b/agent-team/agent_team/db/schema.py @@ -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: diff --git a/agent-team/agent_team/db/transitions.py b/agent-team/agent_team/db/transitions.py new file mode 100644 index 0000000..1fa4635 --- /dev/null +++ b/agent-team/agent_team/db/transitions.py @@ -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) diff --git a/agent-team/agent_team/graph.py b/agent-team/agent_team/graph.py index a724ec5..9983f3a 100644 --- a/agent-team/agent_team/graph.py +++ b/agent-team/agent_team/graph.py @@ -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, diff --git a/agent-team/agent_team/status_page.py b/agent-team/agent_team/status_page.py index 4842949..e5b8a67 100644 --- a/agent-team/agent_team/status_page.py +++ b/agent-team/agent_team/status_page.py @@ -57,14 +57,11 @@ Config (all env, with defaults): from __future__ import annotations -import html -import json import os import sqlite3 from collections import Counter from dataclasses import asdict, dataclass, field from datetime import datetime, timezone -from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from pathlib import Path from typing import Any @@ -73,18 +70,6 @@ from typing import Any _PACKAGE_DIR = Path(__file__).resolve().parent _DEFAULT_DB = _PACKAGE_DIR.parent / "state" / "agent_team.sqlite" -_DEFAULT_HOST = "0.0.0.0" -_DEFAULT_PORT = 8770 - -#