Turn the read-only status page into an auto-updating visual map of the agent-team DAG (INTAKE -> CLARIFY <-> gate -> PLAN <-> REVIEW -> [BUILD -> VERIFY -> DISPATCH] -> DONE), rendered as hand-rolled inline SVG (no CDN/D3 — the R720 is offline/LAN-only). - Stage model (STAGES) with per-node model/agent role labels (Claude/ GPT-4.1/Gemini/DeepSeek/Slack owner) and gated/role-node flags. - Per-stage live state (idle/active/awaiting-human/parked) + count badge, grouped by current_phase; awaiting-human = OPEN pending_question. - GET /api/state JSON sidecar (snapshot_to_dict); inline vanilla-JS poller fetches it every 4s and repaints node states/counts/cards/clock/tooltip in place (no reload, hover/scroll survive). <noscript> meta-refresh fallback retained. - Hover/focus tooltip per node: short_id, description, status, waiting age. - Existing table view kept as a detail section below the map. - Read-only (mode=ro), fail-safe, no secrets; descriptions escaped for both HTML and the JSON-in-script seed (< / > -> \uXXXX), DOM via textContent.
1297 lines
51 KiB
Python
1297 lines
51 KiB
Python
"""LAN-only, READ-ONLY web status dashboard for the agent-team coordinator.
|
|
|
|
The agent-team is a single durable LangGraph pipeline; "agents" here are the
|
|
per-stage model roles, and what an operator actually wants to see is *which task
|
|
is in which phase* and *which tasks are blocked on the human gate*. This module
|
|
serves that view as a tiny self-refreshing HTML page so Adam can glance at the
|
|
queue without SSHing to the R720 and running ``run-team.py list``.
|
|
|
|
Posture (read this before changing anything):
|
|
|
|
* **READ-ONLY.** The page exposes no mutating endpoints. The SQLite ledger is
|
|
opened read-only (``mode=ro`` / ``immutable``-safe URI) and never written. It
|
|
reads the same two durable sources the coordinator owns:
|
|
- ``pending_questions`` (the human-interaction lifecycle ledger) for "who is
|
|
waiting on what" — reused via :mod:`agent_team.ledger`.
|
|
- the LangGraph SQLite checkpoint (``channel_values`` holds ``task``,
|
|
``current_phase``, ``status``) for each task's phase/description — read
|
|
through the same :class:`SqliteSaver` the coordinator uses in
|
|
:func:`agent_team.graph.build_sqlite_checkpointer`.
|
|
* **Sensitive content.** Task descriptions are operator input and may contain
|
|
sensitive detail (repo names, ticket content, remediation context). The page
|
|
renders them in plain text with HTML-escaping but applies NO authentication.
|
|
It must therefore stay **LAN/VPN-only and never public**. The sh-secrev R720
|
|
VM (10.10.60.120, VLAN 60) has no public NIC and sits behind the UniFi
|
|
firewall, so binding ``0.0.0.0`` exposes it to the LAN/VPN only. No secrets,
|
|
tokens, channel refs, or answer bodies are rendered.
|
|
* **Fail-safe.** A missing, locked, or unreadable DB renders a friendly
|
|
"no data" page; the server never crashes on a bad read.
|
|
|
|
The page is a live visual pipeline map: server-rendered inline SVG of the
|
|
agent-team DAG (``INTAKE -> CLARIFY (human gate) -> PLAN <-> REVIEW ->
|
|
[BUILD -> VERIFY -> DISPATCH] -> DONE``), with each stage coloured by live
|
|
state and badged with a task count. Vanilla inline JS polls a JSON sidecar
|
|
endpoint (``GET /api/state``) every few seconds with ``fetch()`` and updates
|
|
the node states, counts, tooltips, and "last updated" clock *in place* — no
|
|
full-page reload, so hover and scroll survive the refresh. A ``<noscript>``
|
|
meta-refresh is the JS-disabled fallback. Everything (SVG + CSS + JS) is inline
|
|
in the served document; nothing is fetched from a CDN, because the R720 VM is
|
|
offline/LAN-only.
|
|
|
|
Why stdlib ``http.server`` and not FastAPI (which is already a dep): the page is
|
|
two read-only GETs (the HTML map and its JSON sidecar). FastAPI/uvicorn buys
|
|
nothing here (no auth, no async I/O) and would couple this LAN dashboard to the
|
|
optional WS1 HTTP-API dependency set. stdlib ``http.server`` keeps it
|
|
zero-new-deps and importable anywhere the ``agent_team`` package is, which also
|
|
keeps :func:`render_html` and :func:`snapshot_to_dict` trivially unit-testable
|
|
without binding a socket.
|
|
|
|
Config (all env, with defaults):
|
|
|
|
* ``AGENT_TEAM_DB`` — path to the SQLite ledger
|
|
(default ``agent-team/state/agent_team.sqlite`` relative to this package).
|
|
* ``AGENT_TEAM_STATUS_HOST`` — bind host (default ``0.0.0.0``; LAN/VPN-only by
|
|
network posture, see above).
|
|
* ``AGENT_TEAM_STATUS_PORT`` — bind port (default ``8770``).
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import html
|
|
import 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
|
|
|
|
# The package root is .../agent-team/agent_team; the default ledger lives at
|
|
# .../agent-team/state/agent_team.sqlite (mirrors run-team.py's _DEFAULT_DB).
|
|
_PACKAGE_DIR = Path(__file__).resolve().parent
|
|
_DEFAULT_DB = _PACKAGE_DIR.parent / "state" / "agent_team.sqlite"
|
|
|
|
_DEFAULT_HOST = "0.0.0.0"
|
|
_DEFAULT_PORT = 8770
|
|
|
|
# <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"}
|
|
_PARKED_STATUSES = {"parked"}
|
|
|
|
|
|
# --- Pipeline map model -------------------------------------------------------
|
|
#
|
|
# The agent-team is one durable LangGraph DAG. These STAGES are the canonical
|
|
# flow we draw, in order, each carrying the model "agent"/role that owns it (per
|
|
# the design: CLARIFY + PLAN = Claude on subscription; REVIEW = GPT-4.1
|
|
# cross_reviewer; SCAN = Gemini scanner; BUILD = DeepSeek fast_coder; the HUMAN
|
|
# GATE = the Slack owner). ``phases`` lists the Phase enum value(s) that map a
|
|
# live task onto that stage; tasks are grouped onto stages by ``current_phase``.
|
|
#
|
|
# VERIFY/DISPATCH are gated/inert in the current deploy (P3 not live) but are
|
|
# drawn so the operator sees the full intended pipeline. SCAN is a role node
|
|
# (Gemini), not a distinct LangGraph phase, so no task ever lands on it — it
|
|
# renders as a labelled-but-idle capability the BUILD/VERIFY stages can call.
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class Stage:
|
|
"""One node in the rendered pipeline map.
|
|
|
|
``key`` is a stable id used by the SVG + JS. ``phases`` are the Phase enum
|
|
*values* whose tasks belong to this stage (empty for pure role nodes like
|
|
the human gate / scan that no ``current_phase`` ever names).
|
|
"""
|
|
|
|
key: str
|
|
label: str
|
|
agent: str # model / role that owns the stage
|
|
phases: tuple[str, ...] = ()
|
|
gated: bool = False # drawn but inert in the current deploy
|
|
role_node: bool = False # capability, not a Phase tasks land on
|
|
|
|
|
|
# Left-to-right pipeline order. Keep in sync with task_model.Phase.
|
|
STAGES: tuple[Stage, ...] = (
|
|
Stage("intake", "Intake", "coordinator", phases=("intake",)),
|
|
Stage("clarify", "Clarify", "Claude (sub)", phases=("clarify",)),
|
|
Stage("gate", "Human Gate", "Slack owner", role_node=True),
|
|
Stage("plan", "Plan", "Claude (sub)", phases=("plan",)),
|
|
Stage("review", "Review", "GPT-4.1 (cross)", phases=("review",)),
|
|
Stage("scan", "Scan", "Gemini", role_node=True),
|
|
Stage("build", "Build", "DeepSeek (fast)", phases=("build",), gated=True),
|
|
Stage("verify", "Verify", "Claude (sub)", phases=("verify",), gated=True),
|
|
Stage("dispatch", "Dispatch", "GitHub PR", gated=True, role_node=True),
|
|
Stage("done", "Done", "—", phases=("done",)),
|
|
)
|
|
|
|
# Phase value -> stage key, derived once from STAGES.
|
|
_PHASE_TO_STAGE: dict[str, str] = {
|
|
phase: stage.key for stage in STAGES for phase in stage.phases
|
|
}
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class TaskView:
|
|
"""A read-only, render-ready view of one pipeline task.
|
|
|
|
Built from the LangGraph checkpoint's ``channel_values`` joined to any open
|
|
pending question for that thread. ``short_id`` is the first 8 chars of the
|
|
thread_id for compact display; ``thread_id`` is retained for completeness.
|
|
"""
|
|
|
|
thread_id: str
|
|
short_id: str
|
|
task: str
|
|
current_phase: str
|
|
status: str
|
|
waiting: bool = False
|
|
waiting_since: str | None = None
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class Snapshot:
|
|
"""Everything :func:`render_html` needs, as plain data (no DB handle).
|
|
|
|
Keeping this pure and serializable is what makes the renderer unit-testable
|
|
without a socket or a live database.
|
|
"""
|
|
|
|
generated_at: str
|
|
db_path: str
|
|
ok: bool = True
|
|
error: str | None = None
|
|
tasks: list[TaskView] = field(default_factory=list)
|
|
question_counts: dict[str, int] = field(default_factory=dict)
|
|
recent_spend: list[dict[str, Any]] = field(default_factory=list)
|
|
spend_total_usd: float | None = None
|
|
budget_available: bool = False
|
|
|
|
@property
|
|
def active_count(self) -> int:
|
|
return sum(1 for t in self.tasks if t.status in _ACTIVE_STATUSES)
|
|
|
|
@property
|
|
def parked_count(self) -> int:
|
|
return sum(1 for t in self.tasks if t.status in _PARKED_STATUSES)
|
|
|
|
@property
|
|
def waiting_tasks(self) -> list[TaskView]:
|
|
return [t for t in self.tasks if t.waiting]
|
|
|
|
@property
|
|
def progressing_tasks(self) -> list[TaskView]:
|
|
return [t for t in self.tasks if not t.waiting]
|
|
|
|
def tasks_for_stage(self, stage_key: str) -> list[TaskView]:
|
|
"""Tasks whose ``current_phase`` maps onto ``stage_key``.
|
|
|
|
Role nodes (human gate / scan / dispatch) own no phase, so they get no
|
|
tasks here even when work is "logically" at the gate — a CLARIFY task
|
|
awaiting an answer stays on the CLARIFY node and is flagged
|
|
awaiting-human there, which is where the operator looks.
|
|
"""
|
|
return [
|
|
t for t in self.tasks if _PHASE_TO_STAGE.get(t.current_phase) == stage_key
|
|
]
|
|
|
|
def stage_state(self, stage: Stage) -> str:
|
|
"""Live state for a stage node: awaiting_human > active > parked > idle.
|
|
|
|
``awaiting_human`` wins so a blocked stage reads as blocked at a glance.
|
|
Role/gated nodes with no tasks fall through to ``idle``.
|
|
"""
|
|
members = self.tasks_for_stage(stage.key)
|
|
if not members:
|
|
return "idle"
|
|
if any(t.waiting for t in members):
|
|
return "awaiting_human"
|
|
if any(t.status in _PARKED_STATUSES for t in members):
|
|
return "parked"
|
|
if any(t.status in _ACTIVE_STATUSES for t in members):
|
|
return "active"
|
|
return "idle"
|
|
|
|
|
|
def _utc_now_iso() -> str:
|
|
"""Current UTC time as a stable, human-readable string."""
|
|
return datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S UTC")
|
|
|
|
|
|
def _short(thread_id: str) -> str:
|
|
"""First 8 chars of a thread_id, for compact display."""
|
|
return thread_id[:8] if thread_id else "(none)"
|
|
|
|
|
|
def _ro_connect(db_path: Path) -> sqlite3.Connection:
|
|
"""Open ``db_path`` strictly READ-ONLY.
|
|
|
|
Uses the SQLite ``file:...?mode=ro`` URI so the connection can never create
|
|
or write the database. ``sqlite3.Row`` row factory matches the reuse readers
|
|
(:mod:`agent_team.ledger`). Raises if the file does not exist (``mode=ro``
|
|
will not create it) — the caller turns that into the friendly no-data page.
|
|
"""
|
|
uri = f"file:{db_path}?mode=ro"
|
|
conn = sqlite3.connect(uri, uri=True, timeout=2.0)
|
|
conn.row_factory = sqlite3.Row
|
|
return conn
|
|
|
|
|
|
def _enumerate_threads(conn: sqlite3.Connection) -> list[str]:
|
|
"""Return the distinct thread_ids the checkpointer has persisted.
|
|
|
|
The LangGraph ``SqliteSaver`` owns the ``checkpoints`` table; we read it
|
|
READ-ONLY. If the table does not exist yet (no task ever ran), returns ``[]``
|
|
rather than raising.
|
|
"""
|
|
try:
|
|
rows = conn.execute(
|
|
"SELECT DISTINCT thread_id FROM checkpoints ORDER BY thread_id"
|
|
).fetchall()
|
|
except sqlite3.OperationalError:
|
|
# No checkpoints table yet (fresh DB) — not an error for a status view.
|
|
return []
|
|
return [r["thread_id"] for r in rows if r["thread_id"]]
|
|
|
|
|
|
def _channel_values_for(saver: Any, thread_id: str) -> dict[str, Any]:
|
|
"""Read the latest checkpoint's ``channel_values`` for one thread.
|
|
|
|
Mirrors the access pattern the coordinator/graph rely on:
|
|
``saver.get({"configurable": {"thread_id": tid}})`` returns the latest
|
|
checkpoint, whose ``channel_values`` carries ``task`` / ``current_phase`` /
|
|
``status`` / ``qa_history``. ``.get`` may return a checkpoint mapping
|
|
directly or ``None`` if the thread has no committed checkpoint; we tolerate
|
|
both and never raise into the page.
|
|
"""
|
|
try:
|
|
checkpoint = saver.get({"configurable": {"thread_id": thread_id}})
|
|
except Exception:
|
|
return {}
|
|
if not checkpoint:
|
|
return {}
|
|
values = checkpoint.get("channel_values") if isinstance(checkpoint, dict) else None
|
|
return values if isinstance(values, dict) else {}
|
|
|
|
|
|
def build_snapshot(db_path: Path | str | None = None) -> Snapshot:
|
|
"""Read the ledger READ-ONLY and assemble a render-ready :class:`Snapshot`.
|
|
|
|
Fail-safe by contract: any DB problem (missing file, locked, missing tables,
|
|
optional langgraph package absent) yields ``ok=False`` with a short message
|
|
rather than raising. The renderer turns that into the friendly no-data page.
|
|
"""
|
|
resolved = Path(db_path) if db_path is not None else _resolve_db_path()
|
|
generated_at = _utc_now_iso()
|
|
|
|
if not resolved.exists():
|
|
return Snapshot(
|
|
generated_at=generated_at,
|
|
db_path=str(resolved),
|
|
ok=False,
|
|
error=f"ledger not found at {resolved} (no tasks have run yet?)",
|
|
)
|
|
|
|
conn: sqlite3.Connection | None = None
|
|
try:
|
|
conn = _ro_connect(resolved)
|
|
question_counts = _read_question_counts(conn)
|
|
open_by_thread = _read_open_questions(conn)
|
|
thread_ids = _enumerate_threads(conn)
|
|
tasks = _read_tasks(resolved, thread_ids, open_by_thread)
|
|
spend, spend_total, budget_available = _read_recent_spend(conn)
|
|
except sqlite3.OperationalError as exc:
|
|
# Locked / busy / unreadable — never crash, just say so.
|
|
return Snapshot(
|
|
generated_at=generated_at,
|
|
db_path=str(resolved),
|
|
ok=False,
|
|
error=f"could not read ledger (busy or unreadable): {exc}",
|
|
)
|
|
except Exception as exc: # defensive: any unexpected read failure
|
|
return Snapshot(
|
|
generated_at=generated_at,
|
|
db_path=str(resolved),
|
|
ok=False,
|
|
error=f"unexpected error reading ledger: {exc}",
|
|
)
|
|
finally:
|
|
if conn is not None:
|
|
conn.close()
|
|
|
|
return Snapshot(
|
|
generated_at=generated_at,
|
|
db_path=str(resolved),
|
|
ok=True,
|
|
tasks=tasks,
|
|
question_counts=question_counts,
|
|
recent_spend=spend,
|
|
spend_total_usd=spend_total,
|
|
budget_available=budget_available,
|
|
)
|
|
|
|
|
|
def snapshot_to_dict(snapshot: Snapshot) -> dict[str, Any]:
|
|
"""Serialize a :class:`Snapshot` into the ``/api/state`` JSON payload.
|
|
|
|
Pure (no I/O) and self-contained so it is unit-testable without a socket.
|
|
The shape is the contract the inline JS poller depends on:
|
|
|
|
* ``ok`` / ``error`` / ``generated_at`` — header + clock + no-data handling.
|
|
* ``stages`` — one entry per pipeline node in draw order, each with its live
|
|
``state``, a task ``count``, and the member ``tasks`` (for the tooltip).
|
|
* ``tasks`` / ``waiting`` / ``question_counts`` / ``summary`` / ``budget`` —
|
|
the detail table + count cards the page also renders.
|
|
|
|
Task descriptions are operator input; they are emitted as raw strings here
|
|
and the *consumer* (the JS, via ``textContent``/JSON parsing) is responsible
|
|
for not interpreting them as markup. ``json.dumps`` itself escapes the
|
|
JSON-string context, and the renderer additionally hardens the inlined seed
|
|
against ``</script>`` / HTML-comment breakouts (see ``_json_for_script``).
|
|
"""
|
|
if not snapshot.ok:
|
|
return {
|
|
"ok": False,
|
|
"error": snapshot.error or "ledger unavailable",
|
|
"generated_at": snapshot.generated_at,
|
|
"db_path": snapshot.db_path,
|
|
"stages": [_stage_payload(s, []) for s in STAGES],
|
|
"tasks": [],
|
|
"waiting": [],
|
|
"question_counts": {},
|
|
"summary": {
|
|
"active": 0,
|
|
"parked": 0,
|
|
"waiting": 0,
|
|
"open_questions": 0,
|
|
"total": 0,
|
|
},
|
|
"budget": {"available": False},
|
|
}
|
|
|
|
stages = [
|
|
_stage_payload(s, snapshot.tasks_for_stage(s.key), snapshot.stage_state(s))
|
|
for s in STAGES
|
|
]
|
|
qc = snapshot.question_counts
|
|
return {
|
|
"ok": True,
|
|
"error": None,
|
|
"generated_at": snapshot.generated_at,
|
|
"db_path": snapshot.db_path,
|
|
"stages": stages,
|
|
"tasks": [_task_payload(t) for t in snapshot.tasks],
|
|
"waiting": [_task_payload(t) for t in snapshot.waiting_tasks],
|
|
"question_counts": qc,
|
|
"summary": {
|
|
"active": snapshot.active_count,
|
|
"parked": snapshot.parked_count,
|
|
"waiting": len(snapshot.waiting_tasks),
|
|
"open_questions": int(qc.get("open", 0)),
|
|
"total": len(snapshot.tasks),
|
|
},
|
|
"budget": {
|
|
"available": snapshot.budget_available,
|
|
"total_usd": snapshot.spend_total_usd,
|
|
"recent": snapshot.recent_spend,
|
|
},
|
|
}
|
|
|
|
|
|
def _task_payload(t: TaskView) -> dict[str, Any]:
|
|
"""JSON-safe view of one task (the tooltip + detail-table source)."""
|
|
return {
|
|
"thread_id": t.thread_id,
|
|
"short_id": t.short_id,
|
|
"task": t.task,
|
|
"current_phase": t.current_phase,
|
|
"status": t.status,
|
|
"waiting": t.waiting,
|
|
"waiting_since": t.waiting_since,
|
|
}
|
|
|
|
|
|
def _stage_payload(
|
|
stage: Stage,
|
|
members: list[TaskView],
|
|
state: str | None = None,
|
|
) -> dict[str, Any]:
|
|
"""JSON-safe view of one pipeline stage node."""
|
|
payload = asdict(stage)
|
|
payload["state"] = state if state is not None else "idle"
|
|
payload["count"] = len(members)
|
|
payload["tasks"] = [_task_payload(t) for t in members]
|
|
return payload
|
|
|
|
|
|
def _resolve_db_path() -> Path:
|
|
"""Resolve the ledger path from ``AGENT_TEAM_DB`` or the package default."""
|
|
env = os.environ.get("AGENT_TEAM_DB")
|
|
return Path(env) if env else _DEFAULT_DB
|
|
|
|
|
|
def _read_question_counts(conn: sqlite3.Connection) -> dict[str, int]:
|
|
"""Count pending_questions rows by status (open/answered/expired/...)."""
|
|
counts: Counter[str] = Counter()
|
|
try:
|
|
rows = conn.execute(
|
|
"SELECT status, COUNT(*) AS n FROM pending_questions GROUP BY status"
|
|
).fetchall()
|
|
except sqlite3.OperationalError:
|
|
return {}
|
|
for r in rows:
|
|
counts[r["status"]] = int(r["n"])
|
|
return dict(counts)
|
|
|
|
|
|
def _read_open_questions(conn: sqlite3.Connection) -> dict[str, str]:
|
|
"""Map thread_id -> earliest ``posted_at`` of an OPEN question on it.
|
|
|
|
An open question means that thread is blocked on the human gate. We keep the
|
|
oldest posted_at so the page can show how long it has been waiting.
|
|
"""
|
|
waiting: dict[str, str] = {}
|
|
try:
|
|
rows = conn.execute(
|
|
"SELECT thread_id, posted_at FROM pending_questions "
|
|
"WHERE status='open' ORDER BY posted_at ASC"
|
|
).fetchall()
|
|
except sqlite3.OperationalError:
|
|
return {}
|
|
for r in rows:
|
|
tid = r["thread_id"]
|
|
if tid and tid not in waiting:
|
|
waiting[tid] = r["posted_at"] or ""
|
|
return waiting
|
|
|
|
|
|
def _read_tasks(
|
|
db_path: Path,
|
|
thread_ids: list[str],
|
|
open_by_thread: dict[str, str],
|
|
) -> list[TaskView]:
|
|
"""Build a :class:`TaskView` per thread from the LangGraph checkpoint.
|
|
|
|
Opens its own READ-ONLY :class:`SqliteSaver` over the same DB file (the saver
|
|
needs its own connection / serde). If the optional langgraph checkpoint
|
|
package is absent, returns ``[]`` so the page degrades to ledger-only data
|
|
instead of crashing.
|
|
"""
|
|
if not thread_ids:
|
|
return []
|
|
|
|
saver_cm = _readonly_saver(db_path)
|
|
if saver_cm is None:
|
|
return []
|
|
|
|
tasks: list[TaskView] = []
|
|
try:
|
|
with saver_cm as saver:
|
|
for tid in thread_ids:
|
|
values = _channel_values_for(saver, tid)
|
|
status = str(values.get("status") or "unknown")
|
|
tasks.append(
|
|
TaskView(
|
|
thread_id=tid,
|
|
short_id=_short(tid),
|
|
task=str(values.get("task") or "(no description)"),
|
|
current_phase=str(values.get("current_phase") or "unknown"),
|
|
status=status,
|
|
waiting=tid in open_by_thread,
|
|
waiting_since=open_by_thread.get(tid) or None,
|
|
)
|
|
)
|
|
except Exception:
|
|
# Saver construction/read failed entirely — degrade to ledger-only.
|
|
return tasks
|
|
# Waiting tasks first, then by phase, for a sensible default ordering.
|
|
tasks.sort(key=lambda t: (not t.waiting, t.current_phase, t.short_id))
|
|
return tasks
|
|
|
|
|
|
def _readonly_saver(db_path: Path) -> Any:
|
|
"""Return an *un-entered* SqliteSaver context manager over a RO connection.
|
|
|
|
Reuses the same :class:`SqliteSaver` + serde the coordinator builds in
|
|
:func:`agent_team.graph.build_checkpoint_serde`, but over a ``mode=ro``
|
|
connection so this dashboard can never write the checkpoint tables. Returns
|
|
``None`` if the optional package is not installed (pre-deploy scaffolding /
|
|
test envs without the checkpoint dep), letting the caller degrade to
|
|
ledger-only data.
|
|
"""
|
|
try:
|
|
from contextlib import closing, contextmanager
|
|
|
|
from langgraph.checkpoint.sqlite import SqliteSaver
|
|
|
|
from agent_team.graph import build_checkpoint_serde
|
|
except Exception:
|
|
return None
|
|
|
|
serde = build_checkpoint_serde()
|
|
|
|
@contextmanager
|
|
def _cm() -> Any:
|
|
uri = f"file:{db_path}?mode=ro"
|
|
with closing(
|
|
sqlite3.connect(uri, uri=True, check_same_thread=False, timeout=2.0)
|
|
) as conn:
|
|
yield SqliteSaver(conn, serde=serde)
|
|
|
|
return _cm()
|
|
|
|
|
|
def _read_recent_spend(
|
|
conn: sqlite3.Connection,
|
|
) -> tuple[list[dict[str, Any]], float | None, bool]:
|
|
"""Read recent budget_ledger spend, if the table is present/readable.
|
|
|
|
Returns ``(recent_rows, total_usd, available)``. ``available`` is ``False``
|
|
when the table is absent (older DB) — the page then simply omits the budget
|
|
panel rather than showing a misleading zero.
|
|
"""
|
|
try:
|
|
rows = conn.execute(
|
|
"SELECT day_bucket, stage, model, usd_cost, recorded_at "
|
|
"FROM budget_ledger ORDER BY recorded_at DESC LIMIT 10"
|
|
).fetchall()
|
|
total_row = conn.execute(
|
|
"SELECT COALESCE(SUM(usd_cost), 0.0) AS total FROM budget_ledger"
|
|
).fetchone()
|
|
except sqlite3.OperationalError:
|
|
return [], None, False
|
|
|
|
recent = [
|
|
{
|
|
"day_bucket": r["day_bucket"],
|
|
"stage": r["stage"],
|
|
"model": r["model"],
|
|
"usd_cost": float(r["usd_cost"] or 0.0),
|
|
"recorded_at": r["recorded_at"],
|
|
}
|
|
for r in rows
|
|
]
|
|
total = float(total_row["total"]) if total_row is not None else 0.0
|
|
return recent, total, True
|
|
|
|
|
|
# --- HTML rendering (PURE: snapshot -> str, no I/O). -------------------------
|
|
|
|
_CSS = """
|
|
:root { color-scheme: light dark; }
|
|
* { box-sizing: border-box; }
|
|
body {
|
|
font-family: -apple-system, BlinkMacSystemFont, "Segoe UI", Roboto, Helvetica,
|
|
Arial, sans-serif;
|
|
margin: 0; padding: 1.5rem; line-height: 1.45;
|
|
background: #0f1115; color: #e6e6e6;
|
|
}
|
|
h1 { font-size: 1.4rem; margin: 0 0 .25rem; }
|
|
h2 { font-size: 1.05rem; margin: 1.5rem 0 .5rem; color: #9fb4d6; }
|
|
.meta { color: #8a93a6; font-size: .85rem; margin-bottom: 1rem; }
|
|
.cards { display: flex; flex-wrap: wrap; gap: .75rem; margin: .5rem 0 1rem; }
|
|
.card {
|
|
background: #1a1e26; border: 1px solid #2a303c; border-radius: 8px;
|
|
padding: .75rem 1rem; min-width: 7rem;
|
|
}
|
|
.card .num { font-size: 1.6rem; font-weight: 600; }
|
|
.card .lbl { color: #8a93a6; font-size: .8rem; text-transform: uppercase;
|
|
letter-spacing: .04em; }
|
|
table { border-collapse: collapse; width: 100%; margin: .25rem 0 1rem;
|
|
font-size: .9rem; }
|
|
th, td { text-align: left; padding: .45rem .6rem; border-bottom: 1px solid #262b35;
|
|
vertical-align: top; }
|
|
th { color: #9fb4d6; font-weight: 600; }
|
|
td.mono, .mono { font-family: ui-monospace, SFMono-Regular, Menlo, monospace; }
|
|
.badge { display: inline-block; padding: .1rem .5rem; border-radius: 999px;
|
|
font-size: .78rem; font-weight: 600; }
|
|
.badge.active { background: #163a2b; color: #6ee7a8; }
|
|
.badge.waiting_human, .badge.waiting { background: #3a3216; color: #f5d76e; }
|
|
.badge.parked { background: #3a1f16; color: #f3a36e; }
|
|
.badge.done { background: #1c2b3a; color: #7fb3f5; }
|
|
.badge.failed { background: #3a1620; color: #f57f9c; }
|
|
.badge.unknown { background: #2a303c; color: #b8c0cf; }
|
|
.desc { max-width: 48rem; }
|
|
.empty { color: #8a93a6; font-style: italic; padding: .5rem 0; }
|
|
.warn { background: #3a1f16; border: 1px solid #5a3322; color: #f3c79e;
|
|
padding: 1rem; border-radius: 8px; }
|
|
footer { color: #6b7385; font-size: .78rem; margin-top: 2rem; }
|
|
|
|
/* --- Pipeline map ------------------------------------------------------- */
|
|
.map-wrap {
|
|
background: #12151c; border: 1px solid #232a36; border-radius: 10px;
|
|
padding: 1rem; overflow-x: auto; position: relative;
|
|
}
|
|
svg.map { display: block; max-width: 100%; height: auto; }
|
|
/* Stage node states. Fill = the box; stroke = its border accent. */
|
|
.node rect.box { fill: #1a1e26; stroke: #2f3645; stroke-width: 1.5;
|
|
transition: fill .25s, stroke .25s; rx: 9; }
|
|
.node text { fill: #e6e6e6; font-size: 13px; font-weight: 600;
|
|
text-anchor: middle; pointer-events: none; }
|
|
.node text.agent { fill: #8a93a6; font-size: 10.5px; font-weight: 500; }
|
|
.node text.gated-tag { fill: #6b7385; font-size: 9px; font-weight: 600;
|
|
letter-spacing: .06em; }
|
|
.node .count { font-size: 11px; font-weight: 700; fill: #0f1115; }
|
|
.node circle.badge { display: none; }
|
|
.node.has-tasks circle.badge { display: inline; }
|
|
/* state -> accent. idle is the default rect above. */
|
|
.node.active rect.box { fill: #163a2b; stroke: #2f8f5e; }
|
|
.node.active circle.badge { fill: #6ee7a8; }
|
|
.node.awaiting_human rect.box { fill: #3a3216; stroke: #b89a2e; }
|
|
.node.awaiting_human circle.badge { fill: #f5d76e; }
|
|
.node.parked rect.box { fill: #3a1f16; stroke: #a85f33; }
|
|
.node.parked circle.badge { fill: #f3a36e; }
|
|
.node.gated rect.box { stroke-dasharray: 5 4; opacity: .82; }
|
|
.node.role rect.box { fill: #181c28; stroke-dasharray: 2 3; }
|
|
.node:focus { outline: none; }
|
|
.node:focus rect.box, .node:hover rect.box { stroke-width: 2.5; }
|
|
.edge { stroke: #3b4456; stroke-width: 1.6; fill: none;
|
|
marker-end: url(#arrow); }
|
|
.edge.loop { stroke: #4a5266; stroke-dasharray: 4 4; }
|
|
.edge-label { fill: #6b7385; font-size: 9.5px; text-anchor: middle; }
|
|
/* Legend */
|
|
.legend { display: flex; flex-wrap: wrap; gap: 1rem; margin: .75rem 0 0;
|
|
font-size: .8rem; color: #9aa3b4; align-items: center; }
|
|
.legend .key { display: inline-flex; align-items: center; gap: .35rem; }
|
|
.legend .swatch { width: 14px; height: 14px; border-radius: 4px;
|
|
border: 1px solid #2f3645; display: inline-block; }
|
|
.legend .swatch.active { background: #163a2b; border-color: #2f8f5e; }
|
|
.legend .swatch.awaiting_human { background: #3a3216; border-color: #b89a2e; }
|
|
.legend .swatch.parked { background: #3a1f16; border-color: #a85f33; }
|
|
.legend .swatch.idle { background: #1a1e26; }
|
|
/* Tooltip (one node, follows hover/focus). */
|
|
#tip {
|
|
position: fixed; z-index: 10; max-width: 24rem; pointer-events: none;
|
|
background: #0b0d12; border: 1px solid #36405a; border-radius: 8px;
|
|
padding: .55rem .7rem; font-size: .8rem; color: #d6dcea;
|
|
box-shadow: 0 8px 24px rgba(0,0,0,.45); display: none; line-height: 1.35;
|
|
}
|
|
#tip.show { display: block; }
|
|
#tip h4 { margin: 0 0 .35rem; font-size: .82rem; color: #9fb4d6; }
|
|
#tip .agent-line { color: #8a93a6; font-size: .74rem; margin-bottom: .35rem; }
|
|
#tip ul { margin: .2rem 0 0; padding-left: 0; list-style: none; }
|
|
#tip li { margin: .3rem 0; border-top: 1px solid #20283a; padding-top: .3rem; }
|
|
#tip li:first-child { border-top: 0; padding-top: 0; }
|
|
#tip .tid { font-family: ui-monospace, SFMono-Regular, Menlo, monospace;
|
|
color: #7fb3f5; }
|
|
#tip .empty-tip { color: #6b7385; font-style: italic; }
|
|
#tip .age { color: #8a93a6; }
|
|
.live-dot { width: .55rem; height: .55rem; border-radius: 50%;
|
|
background: #6ee7a8; display: inline-block; margin-right: .3rem;
|
|
vertical-align: middle; }
|
|
.live-dot.stale { background: #f3a36e; }
|
|
@media (prefers-reduced-motion: reduce) {
|
|
.node rect.box, .edge { transition: none; }
|
|
}
|
|
"""
|
|
|
|
|
|
def _esc(value: object) -> str:
|
|
"""HTML-escape any value (task descriptions are untrusted operator input)."""
|
|
return html.escape(str(value), quote=True)
|
|
|
|
|
|
def _badge(text: str, cls: str | None = None) -> str:
|
|
css = (cls or text or "unknown").strip().lower().replace(" ", "_")
|
|
return f'<span class="badge {_esc(css)}">{_esc(text)}</span>'
|
|
|
|
|
|
def _render_task_rows(tasks: list[TaskView]) -> str:
|
|
if not tasks:
|
|
return '<tr><td colspan="4" class="empty">No tasks.</td></tr>'
|
|
out: list[str] = []
|
|
for t in tasks:
|
|
out.append(
|
|
"<tr>"
|
|
f'<td class="mono">{_esc(t.short_id)}</td>'
|
|
f'<td class="desc">{_esc(t.task)}</td>'
|
|
f"<td>{_badge(t.current_phase)}</td>"
|
|
f"<td>{_badge(t.status)}</td>"
|
|
"</tr>"
|
|
)
|
|
return "".join(out)
|
|
|
|
|
|
def _render_waiting(tasks: list[TaskView]) -> str:
|
|
waiting = [t for t in tasks if t.waiting]
|
|
if not waiting:
|
|
return '<p class="empty">No tasks are waiting on the human gate.</p>'
|
|
rows: list[str] = []
|
|
for t in waiting:
|
|
since = _esc(t.waiting_since) if t.waiting_since else "—"
|
|
rows.append(
|
|
"<tr>"
|
|
f'<td class="mono">{_esc(t.short_id)}</td>'
|
|
f'<td class="desc">{_esc(t.task)}</td>'
|
|
f"<td>{_badge(t.current_phase)}</td>"
|
|
f'<td class="mono">{since}</td>'
|
|
"</tr>"
|
|
)
|
|
return (
|
|
"<table><thead><tr><th>Thread</th><th>Task</th><th>Phase</th>"
|
|
"<th>Open since (UTC)</th></tr></thead><tbody>"
|
|
+ "".join(rows)
|
|
+ "</tbody></table>"
|
|
)
|
|
|
|
|
|
def _render_budget(snapshot: Snapshot) -> str:
|
|
if not snapshot.budget_available:
|
|
return ""
|
|
total = (
|
|
f"${snapshot.spend_total_usd:,.4f}"
|
|
if snapshot.spend_total_usd is not None
|
|
else "—"
|
|
)
|
|
rows: list[str] = []
|
|
for entry in snapshot.recent_spend:
|
|
rows.append(
|
|
"<tr>"
|
|
f'<td class="mono">{_esc(entry.get("recorded_at"))}</td>'
|
|
f"<td>{_esc(entry.get('stage') or '—')}</td>"
|
|
f'<td class="mono">{_esc(entry.get("model"))}</td>'
|
|
f'<td class="mono">${float(entry.get("usd_cost") or 0.0):,.4f}</td>'
|
|
"</tr>"
|
|
)
|
|
body = (
|
|
"".join(rows)
|
|
if rows
|
|
else '<tr><td colspan="4" class="empty">No spend recorded.</td></tr>'
|
|
)
|
|
return (
|
|
f"<h2>Budget — total recorded spend {total}</h2>"
|
|
"<table><thead><tr><th>Recorded (UTC)</th><th>Stage</th>"
|
|
"<th>Model</th><th>USD</th></tr></thead><tbody>" + body + "</tbody></table>"
|
|
)
|
|
|
|
|
|
# --- Inline SVG pipeline map (hand-rolled, no CDN, no D3). -------------------
|
|
#
|
|
# Layout: a serpentine of two rows of five stage nodes, flowing row 1 L->R then
|
|
# row 2 L->R (with a connector down the right side), so the whole DAG reads in
|
|
# one screen. Loop-backs (CLARIFY <-> human gate, PLAN <-> REVIEW) are drawn as
|
|
# dashed curved edges so the operator can see the gating cycles.
|
|
|
|
_NODE_W = 132
|
|
_NODE_H = 66
|
|
_COL_GAP = 56
|
|
_ROW_GAP = 86
|
|
_MARGIN = 24
|
|
_COLS = 5
|
|
|
|
|
|
def _node_xy(index: int) -> tuple[int, int]:
|
|
"""Top-left (x, y) of the stage box at ``index`` in the serpentine grid."""
|
|
row, col = divmod(index, _COLS)
|
|
x = _MARGIN + col * (_NODE_W + _COL_GAP)
|
|
y = _MARGIN + row * (_NODE_H + _ROW_GAP)
|
|
return x, y
|
|
|
|
|
|
def _render_node(index: int, stage: Stage, state: str, count: int) -> str:
|
|
"""One stage box: state class, label, agent role, gated tag, count badge."""
|
|
x, y = _node_xy(index)
|
|
classes = ["node", _esc(state)]
|
|
if stage.gated:
|
|
classes.append("gated")
|
|
if stage.role_node:
|
|
classes.append("role")
|
|
if count > 0:
|
|
classes.append("has-tasks")
|
|
cls = " ".join(classes)
|
|
cx, cy = x + _NODE_W, y # badge anchor: top-right corner
|
|
label_x = x + _NODE_W // 2
|
|
gated_tag = (
|
|
f'<text class="gated-tag" x="{label_x}" y="{y + _NODE_H - 8}">GATED</text>'
|
|
if stage.gated
|
|
else ""
|
|
)
|
|
badge = (
|
|
f'<circle class="badge" cx="{cx}" cy="{cy}" r="11"></circle>'
|
|
f'<text class="count" x="{cx}" y="{cy + 4}">{count}</text>'
|
|
)
|
|
return (
|
|
f'<g class="{cls}" tabindex="0" role="button" '
|
|
f'data-stage="{_esc(stage.key)}" '
|
|
f'aria-label="{_esc(stage.label)} stage, {_esc(state)}, '
|
|
f'{count} task(s)">'
|
|
f'<rect class="box" x="{x}" y="{y}" width="{_NODE_W}" '
|
|
f'height="{_NODE_H}" rx="9"></rect>'
|
|
f'<text x="{label_x}" y="{y + 27}">{_esc(stage.label)}</text>'
|
|
f'<text class="agent" x="{label_x}" y="{y + 45}">'
|
|
f"{_esc(stage.agent)}</text>"
|
|
f"{gated_tag}{badge}</g>"
|
|
)
|
|
|
|
|
|
def _straight_edge(a: int, b: int) -> str:
|
|
"""Forward edge from node ``a`` to node ``b`` (adjacent in the spine)."""
|
|
ax, ay = _node_xy(a)
|
|
bx, by = _node_xy(b)
|
|
a_row = a // _COLS
|
|
b_row = b // _COLS
|
|
if a_row == b_row:
|
|
# Same row: right edge of a -> left edge of b.
|
|
x1, y1 = ax + _NODE_W, ay + _NODE_H // 2
|
|
x2, y2 = bx, by + _NODE_H // 2
|
|
mx = (x1 + x2) / 2
|
|
return (
|
|
f'<path class="edge" d="M{x1},{y1} C{mx},{y1} {mx},{y2} {x2},{y2}"></path>'
|
|
)
|
|
# Row wrap: drop from bottom of a, curve down to top of b.
|
|
x1, y1 = ax + _NODE_W // 2, ay + _NODE_H
|
|
x2, y2 = bx + _NODE_W // 2, by
|
|
my = (y1 + y2) / 2
|
|
return f'<path class="edge" d="M{x1},{y1} C{x1},{my} {x2},{my} {x2},{y2}"></path>'
|
|
|
|
|
|
def _loop_edge(a: int, b: int, label: str) -> str:
|
|
"""Dashed loop-back edge a<->b drawn as an arc above the two same-row nodes."""
|
|
ax, ay = _node_xy(a)
|
|
bx, _ = _node_xy(b)
|
|
# Arc over the top of both boxes (a is to the right of b on the spine).
|
|
x1, y1 = ax + _NODE_W // 2, ay
|
|
x2, y2 = bx + _NODE_W // 2, ay
|
|
lift = 30
|
|
mid_x = (x1 + x2) / 2
|
|
path = (
|
|
f'<path class="edge loop" d="M{x1},{y1} '
|
|
f'C{x1},{y1 - lift} {x2},{y2 - lift} {x2},{y2}"></path>'
|
|
)
|
|
lbl = (
|
|
f'<text class="edge-label" x="{mid_x}" y="{y1 - lift + 2}">{_esc(label)}</text>'
|
|
)
|
|
return path + lbl
|
|
|
|
|
|
def _render_map_svg(snapshot: Snapshot) -> str:
|
|
"""Hand-rolled inline SVG of the agent-team DAG, with live per-stage state.
|
|
|
|
Deterministic geometry (no measurement), all inline — no external assets.
|
|
The JS poller re-paints node classes/counts in place; this server render is
|
|
the first paint and the no-JS view.
|
|
"""
|
|
rows = (len(STAGES) + _COLS - 1) // _COLS
|
|
width = _MARGIN * 2 + _COLS * _NODE_W + (_COLS - 1) * _COL_GAP
|
|
height = _MARGIN * 2 + rows * _NODE_H + (rows - 1) * _ROW_GAP
|
|
|
|
parts: list[str] = []
|
|
# Forward spine: 0->1->...->n following serpentine order.
|
|
for i in range(len(STAGES) - 1):
|
|
parts.append(_straight_edge(i, i + 1))
|
|
# Loop-backs by stage key (robust to index shifts).
|
|
idx = {s.key: i for i, s in enumerate(STAGES)}
|
|
if "clarify" in idx and "gate" in idx:
|
|
parts.append(_loop_edge(idx["gate"], idx["clarify"], "answer ↺"))
|
|
if "plan" in idx and "review" in idx:
|
|
parts.append(_loop_edge(idx["review"], idx["plan"], "revise ↺"))
|
|
|
|
for i, stage in enumerate(STAGES):
|
|
members = snapshot.tasks_for_stage(stage.key)
|
|
state = snapshot.stage_state(stage)
|
|
parts.append(_render_node(i, stage, state, len(members)))
|
|
|
|
arrow = (
|
|
'<defs><marker id="arrow" viewBox="0 0 10 10" refX="9" refY="5" '
|
|
'markerWidth="7" markerHeight="7" orient="auto-start-reverse">'
|
|
'<path d="M0,0 L10,5 L0,10 z" fill="#3b4456"></path></marker></defs>'
|
|
)
|
|
return (
|
|
f'<svg class="map" viewBox="0 0 {width} {height}" '
|
|
f'width="{width}" height="{height}" role="img" '
|
|
'aria-label="agent-team pipeline map">'
|
|
f"{arrow}{''.join(parts)}</svg>"
|
|
)
|
|
|
|
|
|
def _render_legend() -> str:
|
|
return (
|
|
'<div class="legend">'
|
|
'<span class="key"><span class="swatch idle"></span>idle</span>'
|
|
'<span class="key"><span class="swatch active"></span>active</span>'
|
|
'<span class="key"><span class="swatch awaiting_human"></span>'
|
|
"awaiting human</span>"
|
|
'<span class="key"><span class="swatch parked"></span>parked</span>'
|
|
'<span class="key" style="color:#6b7385">'
|
|
"dashed border = gated / not-yet-live stage</span>"
|
|
"</div>"
|
|
)
|
|
|
|
|
|
def _json_for_script(payload: dict[str, Any]) -> str:
|
|
"""JSON safe to inline inside a ``<script>`` element.
|
|
|
|
``json.dumps`` handles the JSON-string escaping. We then escape ``<`` and
|
|
``>`` to their ``\\uXXXX`` JSON escapes, which is the standard hardening for
|
|
JSON embedded in HTML: it makes it impossible for any value (e.g. a task
|
|
description containing ``</script>`` or ``<script>``) to introduce a tag,
|
|
comment, or CDATA breakout, while remaining byte-for-byte parseable by
|
|
``JSON.parse``. ``ensure_ascii=True`` additionally escapes the unicode line
|
|
separators that would otherwise terminate a JS string. Descriptions still
|
|
only ever reach the DOM via ``textContent``, never ``innerHTML``.
|
|
"""
|
|
raw = json.dumps(payload, ensure_ascii=True)
|
|
return raw.replace("<", "\\u003c").replace(">", "\\u003e")
|
|
|
|
|
|
def _poller_script(payload: dict[str, Any]) -> str:
|
|
"""Inline vanilla-JS poller: fetch /api/state, repaint the map in place.
|
|
|
|
No libraries, no CDN. Updates node state classes + count badges, the count
|
|
cards, the "last updated" clock, and the tooltip dataset — without reloading
|
|
the page, so hover/scroll/focus survive. Descriptions reach the DOM only via
|
|
``textContent`` (never innerHTML), so they cannot inject markup.
|
|
"""
|
|
seed = _json_for_script(payload)
|
|
return (
|
|
"<script>(function(){"
|
|
"var POLL=" + str(_POLL_INTERVAL_MS) + ";"
|
|
"var state=" + seed + ";"
|
|
"var tip=document.getElementById('tip');"
|
|
"var clock=document.getElementById('clock');"
|
|
"var dot=document.getElementById('livedot');"
|
|
# ---- helpers ----
|
|
"function age(iso){if(!iso)return '';"
|
|
"var t=Date.parse(iso);if(isNaN(t))return '';"
|
|
"var s=Math.max(0,(Date.now()-t)/1000);"
|
|
"if(s<90)return Math.round(s)+'s ago';"
|
|
"if(s<5400)return Math.round(s/60)+'m ago';"
|
|
"if(s<129600)return Math.round(s/3600)+'h ago';"
|
|
"return Math.round(s/86400)+'d ago';}"
|
|
"function byKey(k){var s=state.stages||[];"
|
|
"for(var i=0;i<s.length;i++){if(s[i].key===k)return s[i];}return null;}"
|
|
"function setNum(id,v){var e=document.getElementById(id);"
|
|
"if(e)e.textContent=String(v==null?0:v);}"
|
|
# ---- repaint the SVG nodes + cards from `state` ----
|
|
"function paint(){"
|
|
"var nodes=document.querySelectorAll('g.node[data-stage]');"
|
|
"for(var i=0;i<nodes.length;i++){var g=nodes[i];"
|
|
"var st=byKey(g.getAttribute('data-stage'));if(!st)continue;"
|
|
"['idle','active','awaiting_human','parked','has-tasks']"
|
|
".forEach(function(c){g.classList.remove(c);});"
|
|
"g.classList.add(st.state||'idle');"
|
|
"if((st.count||0)>0)g.classList.add('has-tasks');"
|
|
"var cnt=g.querySelector('text.count');"
|
|
"if(cnt)cnt.textContent=String(st.count||0);"
|
|
"g.setAttribute('aria-label',(st.label||g.getAttribute('data-stage'))"
|
|
"+' stage, '+(st.state||'idle')+', '+(st.count||0)+' task(s)');}"
|
|
"var sm=state.summary||{};"
|
|
"setNum('card-active',sm.active);setNum('card-parked',sm.parked);"
|
|
"setNum('card-waiting',sm.waiting);setNum('card-open',sm.open_questions);"
|
|
"setNum('card-total',sm.total);}"
|
|
# ---- tooltip ----
|
|
"function fillTip(st){"
|
|
"var h=document.createElement('h4');h.textContent=(st.label||'')+' — '"
|
|
"+(st.state||'idle').replace('_',' ');"
|
|
"var a=document.createElement('div');a.className='agent-line';"
|
|
"a.textContent='agent: '+(st.agent||'—');"
|
|
"tip.innerHTML='';tip.appendChild(h);tip.appendChild(a);"
|
|
"var tasks=st.tasks||[];"
|
|
"if(!tasks.length){var e=document.createElement('div');"
|
|
"e.className='empty-tip';"
|
|
"e.textContent=(st.gated?'gated / not yet live — ':'')+'no tasks here';"
|
|
"tip.appendChild(e);return;}"
|
|
"var ul=document.createElement('ul');"
|
|
"for(var i=0;i<tasks.length;i++){var t=tasks[i];"
|
|
"var li=document.createElement('li');"
|
|
"var idline=document.createElement('div');"
|
|
"var sp=document.createElement('span');sp.className='tid';"
|
|
"sp.textContent=t.short_id||t.thread_id;"
|
|
"idline.appendChild(sp);"
|
|
"idline.appendChild(document.createTextNode(' · '+(t.status||'')));"
|
|
"li.appendChild(idline);"
|
|
"var d=document.createElement('div');d.textContent=t.task||'(no description)';"
|
|
"li.appendChild(d);"
|
|
"var when=t.waiting?t.waiting_since:null;"
|
|
"if(when){var ag=document.createElement('div');ag.className='age';"
|
|
"ag.textContent='waiting since '+when+(age(when)?' ('+age(when)+')':'');"
|
|
"li.appendChild(ag);}"
|
|
"ul.appendChild(li);}"
|
|
"tip.appendChild(ul);}"
|
|
"function showTip(key,ev){var st=byKey(key);if(!st)return;"
|
|
"fillTip(st);tip.classList.add('show');moveTip(ev);}"
|
|
"function moveTip(ev){if(!ev)return;var pad=14;"
|
|
"var x=ev.clientX+pad,y=ev.clientY+pad;"
|
|
"var r=tip.getBoundingClientRect();"
|
|
"if(x+r.width>window.innerWidth)x=ev.clientX-r.width-pad;"
|
|
"if(y+r.height>window.innerHeight)y=ev.clientY-r.height-pad;"
|
|
"tip.style.left=x+'px';tip.style.top=y+'px';}"
|
|
"function hideTip(){tip.classList.remove('show');}"
|
|
# ---- wire node hover/focus ----
|
|
"function wire(){var nodes=document.querySelectorAll('g.node[data-stage]');"
|
|
"for(var i=0;i<nodes.length;i++){(function(g){"
|
|
"var k=g.getAttribute('data-stage');"
|
|
"g.addEventListener('mouseenter',function(e){showTip(k,e);});"
|
|
"g.addEventListener('mousemove',function(e){moveTip(e);});"
|
|
"g.addEventListener('mouseleave',hideTip);"
|
|
"g.addEventListener('focus',function(){var st=byKey(k);if(st){"
|
|
"fillTip(st);var b=g.getBoundingClientRect();tip.classList.add('show');"
|
|
"tip.style.left=(b.left)+'px';tip.style.top=(b.bottom+8)+'px';}});"
|
|
"g.addEventListener('blur',hideTip);"
|
|
"})(nodes[i]);}}"
|
|
# ---- clock ----
|
|
"function setClock(){if(clock)clock.textContent=state.generated_at||'';"
|
|
"if(dot){dot.classList.remove('stale');}}"
|
|
# ---- poll ----
|
|
"function apply(s){state=s;paint();setClock();}"
|
|
"function poll(){fetch('/api/state',{cache:'no-store'})"
|
|
".then(function(r){return r.json();})"
|
|
".then(function(s){apply(s);})"
|
|
".catch(function(){if(dot)dot.classList.add('stale');});}"
|
|
# ---- init ----
|
|
"paint();wire();setClock();"
|
|
"setInterval(poll,POLL);"
|
|
"})();</script>"
|
|
)
|
|
|
|
|
|
def render_html(snapshot: Snapshot) -> str:
|
|
"""Render a :class:`Snapshot` to a complete, self-contained HTML document.
|
|
|
|
Pure function: no I/O, no DB, no socket — unit-testable against a fake
|
|
snapshot. Inline CSS + inline SVG map + inline JS poller only, plus a
|
|
``<noscript>`` meta-refresh fallback, so the page works fully offline (no
|
|
external CDN) and updates live without a full reload.
|
|
"""
|
|
payload = snapshot_to_dict(snapshot)
|
|
head = (
|
|
'<!doctype html><html lang="en"><head><meta charset="utf-8">'
|
|
'<meta name="viewport" content="width=device-width, initial-scale=1">'
|
|
# No-JS fallback: full reload. The inline poller refreshes live without
|
|
# reloading, so this only fires when JS is disabled.
|
|
"<noscript>"
|
|
f'<meta http-equiv="refresh" content="{_REFRESH_SECONDS}">'
|
|
"</noscript>"
|
|
"<title>agent-team status</title>"
|
|
f"<style>{_CSS}</style></head><body>"
|
|
)
|
|
header = (
|
|
"<h1>Sea Haven agent-team — coordinator status</h1>"
|
|
'<div class="meta"><span class="live-dot" id="livedot"></span>'
|
|
'last updated <span class="mono" id="clock">'
|
|
f"{_esc(snapshot.generated_at)}</span> · live "
|
|
f"(polls /api/state every {_POLL_INTERVAL_MS // 1000}s) · "
|
|
f'<span class="mono">{_esc(snapshot.db_path)}</span> · '
|
|
"read-only, LAN/VPN-only</div>"
|
|
)
|
|
tooltip = '<div id="tip" role="tooltip"></div>'
|
|
|
|
# The map + tooltip + poller render on every path (including no-data) so the
|
|
# page goes live the moment the coordinator produces state, without a reload.
|
|
map_section = (
|
|
"<h2>Pipeline map</h2>"
|
|
'<div class="map-wrap">'
|
|
+ _render_map_svg(snapshot)
|
|
+ "</div>"
|
|
+ _render_legend()
|
|
)
|
|
|
|
if not snapshot.ok:
|
|
warn = (
|
|
'<div class="warn"><strong>No data.</strong><br>'
|
|
f"{_esc(snapshot.error or 'ledger unavailable')}<br>"
|
|
"The dashboard is read-only and waits for the coordinator to "
|
|
"produce state. The map will populate live once it does.</div>"
|
|
)
|
|
return (
|
|
head
|
|
+ header
|
|
+ tooltip
|
|
+ map_section
|
|
+ warn
|
|
+ _footer()
|
|
+ _poller_script(payload)
|
|
+ "</body></html>"
|
|
)
|
|
|
|
qc = snapshot.question_counts
|
|
cards = (
|
|
'<div class="cards">'
|
|
f'<div class="card"><div class="num" id="card-active">'
|
|
f"{snapshot.active_count}</div>"
|
|
'<div class="lbl">Active</div></div>'
|
|
f'<div class="card"><div class="num" id="card-parked">'
|
|
f"{snapshot.parked_count}</div>"
|
|
'<div class="lbl">Parked</div></div>'
|
|
f'<div class="card"><div class="num" id="card-waiting">'
|
|
f"{len(snapshot.waiting_tasks)}</div>"
|
|
'<div class="lbl">Waiting (gate)</div></div>'
|
|
f'<div class="card"><div class="num" id="card-open">'
|
|
f"{qc.get('open', 0)}</div>"
|
|
'<div class="lbl">Open questions</div></div>'
|
|
f'<div class="card"><div class="num" id="card-total">'
|
|
f"{len(snapshot.tasks)}</div>"
|
|
'<div class="lbl">Total tasks</div></div>'
|
|
"</div>"
|
|
)
|
|
|
|
waiting_section = "<h2>Waiting on the human gate</h2>" + _render_waiting(
|
|
snapshot.tasks
|
|
)
|
|
|
|
tasks_section = (
|
|
"<h2>All tasks</h2>"
|
|
"<table><thead><tr><th>Thread</th><th>Task</th><th>Phase</th>"
|
|
"<th>Status</th></tr></thead><tbody>"
|
|
+ _render_task_rows(snapshot.tasks)
|
|
+ "</tbody></table>"
|
|
)
|
|
|
|
budget_section = _render_budget(snapshot)
|
|
|
|
return (
|
|
head
|
|
+ header
|
|
+ tooltip
|
|
+ cards
|
|
+ map_section
|
|
+ waiting_section
|
|
+ tasks_section
|
|
+ budget_section
|
|
+ _footer()
|
|
+ _poller_script(payload)
|
|
+ "</body></html>"
|
|
)
|
|
|
|
|
|
def _footer() -> str:
|
|
return (
|
|
"<footer>READ-ONLY status dashboard · no mutating endpoints "
|
|
"· unauthenticated — do not expose to the public "
|
|
"internet. Task descriptions may be sensitive.</footer>"
|
|
)
|
|
|
|
|
|
# --- HTTP server (thin shell around the pure renderer). ----------------------
|
|
|
|
|
|
def _make_handler(db_path: Path) -> type[BaseHTTPRequestHandler]:
|
|
"""Build a request handler bound to ``db_path`` (closure, no globals)."""
|
|
|
|
class _StatusHandler(BaseHTTPRequestHandler):
|
|
server_version = "AgentTeamStatus/1.0"
|
|
|
|
# Strip any query string so "/api/state?ts=123" (a cache-buster the JS
|
|
# may append) still routes. Read-only: only GET, only these two routes.
|
|
def _route(self) -> str:
|
|
return self.path.split("?", 1)[0]
|
|
|
|
def do_GET(self) -> None: # noqa: N802 (BaseHTTPRequestHandler API)
|
|
route = self._route()
|
|
# Favicon gets a 204 so browsers stop logging 404s.
|
|
if route == "/favicon.ico":
|
|
self.send_response(204)
|
|
self.end_headers()
|
|
return
|
|
if route == "/api/state":
|
|
self._serve_json()
|
|
return
|
|
# Every other path serves the dashboard HTML.
|
|
self._serve_html()
|
|
|
|
def _serve_json(self) -> None:
|
|
try:
|
|
payload = snapshot_to_dict(build_snapshot(db_path))
|
|
except Exception as exc: # last-ditch: never 500 the sidecar
|
|
payload = {
|
|
"ok": False,
|
|
"error": f"snapshot error: {exc}",
|
|
"generated_at": _utc_now_iso(),
|
|
"db_path": str(db_path),
|
|
"stages": [],
|
|
"tasks": [],
|
|
"waiting": [],
|
|
"question_counts": {},
|
|
"summary": {},
|
|
"budget": {"available": False},
|
|
}
|
|
body = json.dumps(payload).encode("utf-8")
|
|
self.send_response(200)
|
|
self.send_header("Content-Type", "application/json; charset=utf-8")
|
|
self.send_header("Content-Length", str(len(body)))
|
|
self.send_header("X-Content-Type-Options", "nosniff")
|
|
self.send_header("Referrer-Policy", "no-referrer")
|
|
self.send_header("Cache-Control", "no-store")
|
|
self.end_headers()
|
|
self.wfile.write(body)
|
|
|
|
def _serve_html(self) -> None:
|
|
try:
|
|
snapshot = build_snapshot(db_path)
|
|
page = render_html(snapshot)
|
|
except Exception as exc: # last-ditch: never 500 the dashboard
|
|
page = render_html(
|
|
Snapshot(
|
|
generated_at=_utc_now_iso(),
|
|
db_path=str(db_path),
|
|
ok=False,
|
|
error=f"render error: {exc}",
|
|
)
|
|
)
|
|
body = page.encode("utf-8")
|
|
self.send_response(200)
|
|
self.send_header("Content-Type", "text/html; charset=utf-8")
|
|
self.send_header("Content-Length", str(len(body)))
|
|
# Defense-in-depth headers for an unauthenticated LAN page.
|
|
self.send_header("X-Content-Type-Options", "nosniff")
|
|
self.send_header("Referrer-Policy", "no-referrer")
|
|
self.send_header("Cache-Control", "no-store")
|
|
self.end_headers()
|
|
self.wfile.write(body)
|
|
|
|
def log_message(self, fmt: str, *args: Any) -> None:
|
|
# Quiet by default; access logs are noise for a refresh-loop page.
|
|
return
|
|
|
|
return _StatusHandler
|
|
|
|
|
|
def serve(
|
|
host: str | None = None,
|
|
port: int | None = None,
|
|
db_path: Path | str | None = None,
|
|
) -> None:
|
|
"""Start the LAN-only read-only status server (blocking).
|
|
|
|
Resolves host/port/db from args or env (``AGENT_TEAM_STATUS_HOST`` default
|
|
``0.0.0.0``, ``AGENT_TEAM_STATUS_PORT`` default ``8770``, ``AGENT_TEAM_DB``
|
|
default the package ledger). Binds and serves until interrupted. The server
|
|
is read-only and unauthenticated — see the module docstring on network
|
|
posture (LAN/VPN-only, never public).
|
|
"""
|
|
resolved_host = host or os.environ.get("AGENT_TEAM_STATUS_HOST") or _DEFAULT_HOST
|
|
if port is not None:
|
|
resolved_port = int(port)
|
|
else:
|
|
resolved_port = int(os.environ.get("AGENT_TEAM_STATUS_PORT", _DEFAULT_PORT))
|
|
resolved_db = Path(db_path) if db_path is not None else _resolve_db_path()
|
|
|
|
handler = _make_handler(resolved_db)
|
|
httpd = ThreadingHTTPServer((resolved_host, resolved_port), handler)
|
|
print(
|
|
f"[status_page] serving READ-ONLY dashboard on "
|
|
f"http://{resolved_host}:{resolved_port} (db={resolved_db}); "
|
|
"LAN/VPN-only, unauthenticated."
|
|
)
|
|
try:
|
|
httpd.serve_forever()
|
|
except KeyboardInterrupt:
|
|
pass
|
|
finally:
|
|
httpd.server_close()
|
|
|
|
|
|
if __name__ == "__main__": # pragma: no cover - manual / systemd entrypoint
|
|
serve()
|