This repository has been archived on 2026-08-04. You can view files and clone it, but cannot push or open issues or pull requests.
orchestrator/agent-team/agent_team/api.py
Adam Moussa 366d07a84e fix(ws1): skip fastapi TestClient tests when fastapi absent + ruff format
CI has no fastapi (box-only dependency); api.py imports it lazily. Guard the
7 TestClient smoke tests with skipif(find_spec('fastapi') is None) so they
skip in CI instead of failing collection, leaving the 14 invoker_multi tests
running. Also apply ruff format to the 5 WS1 files CI flagged.
2026-06-23 11:40:46 -04:00

303 lines
11 KiB
Python

"""FastAPI HTTP API for the Sea Haven agent-team orchestrator (WS1 — code only).
Exposes two endpoints under bearer-token auth (env ``AGENT_TEAM_API_TOKEN``):
``POST /tasks``
Intake a new agent-team task: creates a task via :meth:`Coordinator.start_task`
and returns the minted ``thread_id``. Requires ``task`` in the JSON body.
``GET /tasks/{thread_id}``
Return the current pipeline state for a running task from the LangGraph
checkpoint. Returns ``{"thread_id": ..., "state": {...}}`` when found, 404
when the thread has no checkpoint.
``POST /orchestrator/invoke``
Route a one-shot prompt through the orchestrator's root ``run.py`` (the
existing multi-model router). Returns ``{"text": ...}`` with the model's
response. Accepts ``{"prompt": "..."}`` in the JSON body.
Security
--------
* **Bearer-token auth**: every request must carry ``Authorization: Bearer <token>``
where the token matches ``AGENT_TEAM_API_TOKEN``. Missing or wrong token → 401.
An empty / absent ``AGENT_TEAM_API_TOKEN`` at startup raises :class:`RuntimeError`
so the server refuses to start in an unconfigured state (no open-by-default).
* **Bind to 127.0.0.1** by default: :func:`make_app` accepts a ``host`` kwarg
defaulting to ``"127.0.0.1"`` so a misconfigured caller cannot accidentally
bind to 0.0.0.0. The uvicorn run call at the bottom of :func:`serve` enforces
this.
* **Token comparison is constant-time** (``hmac.compare_digest``) to prevent
timing side-channels on the secret.
* **No secrets in code**: the token is read from env only, never hardcoded here.
Code-only gate (do NOT launch)
------------------------------
This module is WS1 code scaffolding. Do NOT import-and-run it on the R720 box
until the GPT-4.1 cross-review and /sh-security-review gates clear. The
``serve()`` helper is provided for the attended deploy step only.
"""
from __future__ import annotations
import hmac
import os
import subprocess
import sys
from pathlib import Path
from typing import Any
from pydantic import BaseModel
__all__ = [
"AGENT_TEAM_API_TOKEN_ENV",
"make_app",
"serve",
]
AGENT_TEAM_API_TOKEN_ENV = "AGENT_TEAM_API_TOKEN"
_DEFAULT_HOST = "127.0.0.1"
_DEFAULT_PORT = 8765
# ---------------------------------------------------------------------------
# Pydantic request/response models (module-level so FastAPI resolves them).
# ---------------------------------------------------------------------------
class TaskRequest(BaseModel):
task: str
transport: str = "claude_code"
class TaskResponse(BaseModel):
thread_id: str
class TaskStateResponse(BaseModel):
thread_id: str
state: dict[str, Any]
class OrchestratorInvokeRequest(BaseModel):
prompt: str
class OrchestratorInvokeResponse(BaseModel):
text: str
# ---------------------------------------------------------------------------
# Token helper.
# ---------------------------------------------------------------------------
def _get_token() -> str:
"""Read the API bearer token from the environment.
Raises :class:`RuntimeError` if the variable is absent or empty — the server
must never start without a configured secret.
"""
token = os.environ.get(AGENT_TEAM_API_TOKEN_ENV, "").strip()
if not token:
raise RuntimeError(
f"{AGENT_TEAM_API_TOKEN_ENV} is not set or empty; refusing to start "
"the HTTP API without a configured bearer token."
)
return token
def _resolve_run_py() -> Path:
"""Resolve the orchestrator root's ``run.py``.
``api.py`` lives at ``<root>/agent-team/agent_team/api.py``, so the root is
``parents[2]``.
"""
return Path(__file__).resolve().parents[2] / "run.py"
def _build_coordinator(db_path: str | Path | None = None) -> Any:
"""Build a default :class:`Coordinator` for the HTTP API.
Uses a null transport by default because the HTTP API is not directly wired
to a Slack/GitHub channel — the caller provides answers through the ledger.
The ledger path defaults to the same ``state/agent_team.sqlite`` the CLI uses.
"""
from agent_team.coordinator import Coordinator # noqa: PLC0415
from agent_team.transport.base import Transport # noqa: PLC0415
class _NullTransport(Transport):
def post_question(self, **kw: Any) -> str: # type: ignore[override]
return f"api:{kw.get('question_id', 'unknown')}"
def parse_answer(self, raw: Any) -> tuple[str, Any, str]: # type: ignore[override]
raise NotImplementedError("HTTP API does not use transport.parse_answer")
if db_path is None:
db_path = Path(__file__).resolve().parent.parent / "state" / "agent_team.sqlite"
return Coordinator(db_path=db_path, transport=_NullTransport())
# ---------------------------------------------------------------------------
# App factory.
# ---------------------------------------------------------------------------
def make_app(
*,
coordinator: Any = None,
db_path: str | Path | None = None,
host: str = _DEFAULT_HOST,
) -> Any:
"""Build and return the FastAPI application (does NOT start it).
``coordinator`` is an optional pre-built :class:`~agent_team.coordinator.Coordinator`
instance. When omitted, :func:`_build_coordinator` constructs one from
``db_path`` at request time.
``host`` defaults to ``"127.0.0.1"`` — do NOT change to ``"0.0.0.0"`` without
the security gate.
"""
try:
from fastapi import Depends, FastAPI, HTTPException # noqa: PLC0415
from fastapi.security import ( # noqa: PLC0415
HTTPAuthorizationCredentials,
HTTPBearer,
)
except ImportError as exc:
raise RuntimeError(
"fastapi is unavailable; install it to use the HTTP API: "
"pip install fastapi uvicorn"
) from exc
app = FastAPI(
title="Sea Haven agent-team API",
description="Internal HTTP API for the agent-team orchestrator (VPN-only, 127.0.0.1).",
version="0.1.0",
)
_bearer = HTTPBearer(auto_error=False)
def _check_token(
creds: HTTPAuthorizationCredentials | None = Depends(_bearer),
) -> None:
expected = _get_token()
if creds is None or not hmac.compare_digest(creds.credentials, expected):
raise HTTPException(
status_code=401, detail="invalid or missing bearer token"
)
# Coordinator is built lazily and cached on first use so test clients that
# inject a pre-built coordinator never trigger the DB path.
_state: dict[str, Any] = {"coordinator": coordinator, "ready": False}
def _get_coordinator() -> Any:
if _state["coordinator"] is None:
_state["coordinator"] = _build_coordinator(db_path)
coord = _state["coordinator"]
if not _state["ready"]:
coord.setup()
_state["ready"] = True
return coord
@app.post(
"/tasks", response_model=TaskResponse, dependencies=[Depends(_check_token)]
)
def create_task(request: TaskRequest) -> TaskResponse:
"""Intake a new agent-team task and run it to the first human gate."""
coord = _get_coordinator()
thread_id = coord.start_task(
task_text=request.task,
transport_name=request.transport,
)
return TaskResponse(thread_id=thread_id)
@app.get(
"/tasks/{thread_id}",
response_model=TaskStateResponse,
dependencies=[Depends(_check_token)],
)
def get_task(thread_id: str) -> TaskStateResponse:
"""Return the current LangGraph checkpoint state for a task."""
coord = _get_coordinator()
graph = getattr(coord, "graph", None)
if graph is None:
raise HTTPException(status_code=503, detail="coordinator not ready")
try:
snapshot = graph.get_state({"configurable": {"thread_id": thread_id}})
if snapshot is None or not getattr(snapshot, "values", None):
raise HTTPException(
status_code=404, detail=f"no state for thread {thread_id!r}"
)
state_dict: dict[str, Any] = dict(snapshot.values)
except HTTPException:
raise
except Exception as exc:
raise HTTPException(status_code=500, detail=str(exc)) from exc
return TaskStateResponse(thread_id=thread_id, state=state_dict)
@app.post(
"/orchestrator/invoke",
response_model=OrchestratorInvokeResponse,
dependencies=[Depends(_check_token)],
)
def orchestrator_invoke(
request: OrchestratorInvokeRequest,
) -> OrchestratorInvokeResponse:
"""Route a one-shot prompt through the root orchestrator run.py.
Shells out to the orchestrator root's ``run.py`` with the prompt text.
The output is returned raw (UNTRUSTED); callers are responsible for
validating it.
"""
run_py = _resolve_run_py()
if not run_py.exists():
raise HTTPException(
status_code=503,
detail=f"orchestrator run.py not found at {run_py}",
)
try:
completed = subprocess.run( # noqa: S603 - args are not shell-interpolated
[sys.executable, str(run_py), request.prompt],
capture_output=True,
text=True,
check=True,
timeout=600,
)
except subprocess.TimeoutExpired:
raise HTTPException(
status_code=504, detail="orchestrator invocation timed out"
) from None
except subprocess.CalledProcessError as exc:
raise HTTPException(
status_code=500,
detail=f"orchestrator call failed (exit {exc.returncode}): {exc.stderr[:500]}",
) from exc
return OrchestratorInvokeResponse(text=completed.stdout)
return app
# ---------------------------------------------------------------------------
# Attended deploy helper (do NOT call in tests or CI).
# ---------------------------------------------------------------------------
def serve(
*,
host: str = _DEFAULT_HOST,
port: int = _DEFAULT_PORT,
db_path: str | Path | None = None,
) -> None: # pragma: no cover - attended deploy step only
"""Start the uvicorn server (attended deploy — do NOT call in tests or CI).
Binds to ``127.0.0.1`` by default. Change ``host`` only after the
GPT-4.1 cross-review and /sh-security-review gates clear.
"""
try:
import uvicorn # noqa: PLC0415
except ImportError as exc:
raise RuntimeError(
"uvicorn is unavailable; install it: pip install uvicorn"
) from exc
app = make_app(host=host, db_path=db_path)
uvicorn.run(app, host=host, port=port)