Merge remote-tracking branch 'origin/feat/ws1-inprocess-models-http-api' into feat/ws-activation-wiring
This commit is contained in:
commit
e3a10a6260
8 changed files with 1035 additions and 32 deletions
347
agent-team/agent_team/api.py
Normal file
347
agent-team/agent_team/api.py
Normal file
|
|
@ -0,0 +1,347 @@
|
|||
"""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
|
||||
import threading
|
||||
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
|
||||
|
||||
# Bound how many `/orchestrator/invoke` subprocesses can run at once. Each spawns
|
||||
# a 600s-timeout child; without a cap an authenticated caller could exhaust CPU /
|
||||
# memory / file descriptors by firing many concurrent invocations (DoS). The
|
||||
# endpoint runs in Starlette's threadpool, so a threading semaphore is the right
|
||||
# primitive; over-limit requests get 429 rather than queueing unboundedly.
|
||||
_MAX_CONCURRENT_INVOKES = 2
|
||||
_invoke_semaphore = threading.BoundedSemaphore(_MAX_CONCURRENT_INVOKES)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# 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
|
||||
|
||||
# Fail fast at app-build time if the bearer token is unconfigured, so a
|
||||
# misconfigured deploy never reaches a serving state (matches the module
|
||||
# docstring's "refuses to start without a secret" contract). The per-request
|
||||
# check still calls _get_token() so a token rotation/unset after boot is
|
||||
# also caught.
|
||||
_get_token()
|
||||
|
||||
# VPN-only internal API: disable the interactive docs and the OpenAPI schema
|
||||
# so the endpoint surface is not exposed unauthenticated (these routes carry
|
||||
# no auth dependency by design in FastAPI).
|
||||
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",
|
||||
docs_url=None,
|
||||
redoc_url=None,
|
||||
openapi_url=None,
|
||||
)
|
||||
|
||||
_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:
|
||||
# Do not leak the raw exception text (may carry internal paths /
|
||||
# state) into the response body; log it server-side and return a
|
||||
# generic message.
|
||||
print(f"[api] get_task error for {thread_id!r}: {exc!r}", file=sys.stderr)
|
||||
raise HTTPException(
|
||||
status_code=500, detail="internal error retrieving task state"
|
||||
) 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}",
|
||||
)
|
||||
# Bound concurrent subprocess fan-out (DoS guard). Non-blocking acquire:
|
||||
# over the cap we reject with 429 rather than pile up 600s children.
|
||||
if not _invoke_semaphore.acquire(blocking=False):
|
||||
raise HTTPException(
|
||||
status_code=429,
|
||||
detail="too many concurrent orchestrator invocations; retry later",
|
||||
)
|
||||
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:
|
||||
# Log stderr server-side; do not return it in the response body
|
||||
# (may carry internal paths / orchestrator internals).
|
||||
print(
|
||||
f"[api] orchestrator_invoke failed (exit {exc.returncode}): "
|
||||
f"{(exc.stderr or '')[:1000]}",
|
||||
file=sys.stderr,
|
||||
)
|
||||
raise HTTPException(
|
||||
status_code=500,
|
||||
detail=f"orchestrator call failed (exit {exc.returncode})",
|
||||
) from exc
|
||||
finally:
|
||||
_invoke_semaphore.release()
|
||||
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)
|
||||
170
agent-team/agent_team/invoker_multi.py
Normal file
170
agent-team/agent_team/invoker_multi.py
Normal file
|
|
@ -0,0 +1,170 @@
|
|||
"""In-process invokers for non-Claude models (GPT-4.1, DeepSeek, Gemini) — WS1.
|
||||
|
||||
Mirrors the pattern of :mod:`agent_team.invoker` (which supplies Claude invokers
|
||||
for the billing seam) but targets the **non-Claude** model seams:
|
||||
|
||||
* :func:`agent_team.nodes.review_loop.set_review_invoker` — GPT-4.1
|
||||
(``cross_reviewer``)
|
||||
* :func:`agent_team.nodes.builders_llm.default_build` is replaced in-process
|
||||
by binding :func:`make_fast_coder_invoker` as the builder
|
||||
|
||||
Each factory lazily imports the root orchestrator's ``models`` module (never at
|
||||
module load) so importing this file on a machine where ``models.py`` is absent
|
||||
(CI, Mac dev, tests) does not raise :class:`ModuleNotFoundError`. The ``models``
|
||||
module itself is kept out of the agent-team package graph; we reach it at runtime
|
||||
through the orchestrator root's ``sys.path`` entry that ``run-team.py`` and the
|
||||
graph module already bootstrap.
|
||||
|
||||
Model keys accepted by :func:`multi_invoke`:
|
||||
|
||||
``cross_reviewer``
|
||||
GPT-4.1 via ``models.get_cross_reviewer()``. The §3.2 cross-family reviewer.
|
||||
``fast_coder``
|
||||
DeepSeek via ``models.get_fast_coder()``. The P3 mechanical-edit builder.
|
||||
``scanner``
|
||||
Gemini via ``models.get_scanner()``. The P3 security-scan verifier.
|
||||
|
||||
Security
|
||||
--------
|
||||
All model outputs returned by :func:`multi_invoke` are UNTRUSTED. Callers are
|
||||
responsible for defensive parsing before acting on the result. This module
|
||||
performs no parsing — it returns the raw model response text and leaves the
|
||||
fail-safe discipline to the caller (review_loop_llm.parse_verdict,
|
||||
builders_llm._extract_diff, etc.).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import sys
|
||||
from collections.abc import Callable
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
__all__ = [
|
||||
"MODELS",
|
||||
"MultiInvoker",
|
||||
"bind_multi_invoker",
|
||||
"make_fast_coder_invoker",
|
||||
"make_scanner_invoker",
|
||||
"multi_invoke",
|
||||
]
|
||||
|
||||
# The set of model keys this module can dispatch to.
|
||||
MODELS: frozenset[str] = frozenset({"cross_reviewer", "fast_coder", "scanner"})
|
||||
|
||||
# Callable type: (prompt, **kw) -> str — same shape as ReviewInvoker and BuildCallable.
|
||||
MultiInvoker = Callable[..., str]
|
||||
|
||||
|
||||
def _ensure_orchestrator_on_path() -> None:
|
||||
"""Add the orchestrator root to sys.path if absent.
|
||||
|
||||
``models.py`` lives at the orchestrator root (not inside agent-team). This
|
||||
file is at ``<root>/agent-team/agent_team/invoker_multi.py``, so the root
|
||||
is ``parents[2]``. Mirroring the bootstrap in ``run-team.py`` ensures the
|
||||
deferred import in :func:`_load_model` succeeds when ``multi_invoke`` is
|
||||
called in a process that has NOT already bootstrapped ``sys.path``.
|
||||
"""
|
||||
root = str(Path(__file__).resolve().parents[2])
|
||||
if root not in sys.path:
|
||||
sys.path.insert(0, root)
|
||||
|
||||
|
||||
def _load_model(key: str) -> Any:
|
||||
"""Lazily import the orchestrator ``models`` module and return the model.
|
||||
|
||||
``key`` must be one of :data:`MODELS`. Raises :class:`ValueError` on an
|
||||
unknown key and :class:`RuntimeError` when the ``models`` module is
|
||||
unavailable (the orchestrator is not on ``sys.path``).
|
||||
|
||||
Deferred import keeps this module importable everywhere the orchestrator
|
||||
root is absent (CI, tests, Mac dev).
|
||||
"""
|
||||
_ensure_orchestrator_on_path()
|
||||
try:
|
||||
import models # noqa: PLC0415 - intentional deferred import
|
||||
except ImportError as exc:
|
||||
raise RuntimeError(
|
||||
"Cannot import orchestrator 'models' module — ensure the orchestrator "
|
||||
"root is on sys.path (run-team.py bootstraps this automatically). "
|
||||
"Tests should mock multi_invoke rather than call it."
|
||||
) from exc
|
||||
factories = {
|
||||
"cross_reviewer": models.get_cross_reviewer,
|
||||
"fast_coder": models.get_fast_coder,
|
||||
"scanner": models.get_scanner,
|
||||
}
|
||||
if key not in factories:
|
||||
raise ValueError(
|
||||
f"unknown model key {key!r}; must be one of: " + ", ".join(sorted(MODELS))
|
||||
)
|
||||
return factories[key]()
|
||||
|
||||
|
||||
def multi_invoke(prompt: str, *, model: str, **kw: Any) -> str:
|
||||
"""Invoke a non-Claude model in-process and return its response text.
|
||||
|
||||
``model`` must be one of :data:`MODELS`. The underlying model is loaded
|
||||
lazily via :func:`_load_model` so the orchestrator's ``models`` module is
|
||||
not imported until first call. Any call error propagates to the caller, which
|
||||
is responsible for fail-safe handling (e.g. :func:`review_plan` catches all
|
||||
exceptions and returns REQUEST_CHANGES).
|
||||
|
||||
The response text is returned raw (UNTRUSTED); callers must parse/validate
|
||||
before acting on it.
|
||||
"""
|
||||
model_obj = _load_model(model)
|
||||
result = model_obj.invoke(prompt)
|
||||
text = getattr(result, "content", result)
|
||||
return text if isinstance(text, str) else str(text)
|
||||
|
||||
|
||||
def make_fast_coder_invoker() -> MultiInvoker:
|
||||
"""Return an in-process callable that routes prompts to DeepSeek ``fast_coder``.
|
||||
|
||||
Mirrors :func:`agent_team.nodes.review_loop_llm.make_cross_reviewer_invoker`
|
||||
for the builder seam. The callable matches :data:`~agent_team.nodes.builders_llm.BuildCallable`
|
||||
``(instruction: str) -> str`` and is suitable for passing as the ``build``
|
||||
kwarg to :func:`agent_team.nodes.builders_llm.build_candidate_diff`.
|
||||
"""
|
||||
|
||||
def _invoke(instruction: str, **_kw: Any) -> str:
|
||||
return multi_invoke(instruction, model="fast_coder")
|
||||
|
||||
return _invoke
|
||||
|
||||
|
||||
def make_scanner_invoker() -> MultiInvoker:
|
||||
"""Return an in-process callable that routes prompts to Gemini ``scanner``.
|
||||
|
||||
The callable is suitable for injection into the verifier's scan seam.
|
||||
"""
|
||||
|
||||
def _invoke(prompt: str, **_kw: Any) -> str:
|
||||
return multi_invoke(prompt, model="scanner")
|
||||
|
||||
return _invoke
|
||||
|
||||
|
||||
def bind_multi_invoker() -> None:
|
||||
"""Bind in-process non-Claude invokers into the review and builder seams.
|
||||
|
||||
Call once at startup (alongside :func:`agent_team.invoker.bind_subscription_invoker`)
|
||||
to replace the subprocess-based defaults with in-process model calls:
|
||||
|
||||
* ``review_loop.set_review_invoker`` → in-process GPT-4.1 via
|
||||
:func:`~agent_team.nodes.review_loop_llm.make_cross_reviewer_invoker`
|
||||
* Builder default (``builders_llm.default_build``) is already replaced at
|
||||
module load in this build; :func:`make_fast_coder_invoker` is the callable
|
||||
for explicit injection into :func:`~agent_team.nodes.builders_llm.build_candidate_diff`
|
||||
when needed.
|
||||
|
||||
Lazy-imported so importing ``invoker_multi`` has no global side effects and
|
||||
tests that never call ``bind_multi_invoker`` remain model-free.
|
||||
"""
|
||||
from agent_team.nodes import review_loop # noqa: PLC0415
|
||||
from agent_team.nodes.review_loop_llm import ( # noqa: PLC0415
|
||||
make_cross_reviewer_invoker,
|
||||
)
|
||||
|
||||
review_loop.set_review_invoker(make_cross_reviewer_invoker())
|
||||
|
|
@ -61,6 +61,8 @@ __all__ = [
|
|||
"as_diff_builder",
|
||||
"build_candidate_diff",
|
||||
"default_build",
|
||||
"make_fast_coder_invoker",
|
||||
"subprocess_build",
|
||||
]
|
||||
|
||||
# The injectable build seam: given the rendered build instruction (a string),
|
||||
|
|
@ -154,24 +156,41 @@ def _orchestrator_root() -> Path:
|
|||
return Path(__file__).resolve().parents[3]
|
||||
|
||||
|
||||
def default_build(instruction: str, *, route: _OrchestratorRoute | None = None) -> str:
|
||||
"""Default :data:`BuildCallable`: route the build to DeepSeek ``fast_coder``.
|
||||
def make_fast_coder_invoker() -> "BuildCallable":
|
||||
"""Return an in-process :data:`BuildCallable` backed by DeepSeek ``fast_coder``.
|
||||
|
||||
Calls the local orchestrator out-of-process — ``python3 <root>/run.py
|
||||
"<instruction>"`` — and returns its stdout. The orchestrator routes a
|
||||
well-specified coding task to its ``fast_coder`` agent (DeepSeek); this is
|
||||
the design's "DeepSeek mechanical edits, via the local orchestrator" path,
|
||||
deliberately NOT :func:`agent_team.billing.claude_invoke`.
|
||||
Lazy-imports the orchestrator ``models`` module (never at module load) to
|
||||
call ``models.get_fast_coder()`` in-process. Mirrors
|
||||
:func:`~agent_team.nodes.review_loop_llm.make_cross_reviewer_invoker` for
|
||||
the builder seam (WS1). Any error propagates to the caller
|
||||
(:func:`build_candidate_diff`), which wraps errors as a ``no_op`` result.
|
||||
"""
|
||||
|
||||
No orchestrator module is imported at module top (deferred, mirroring
|
||||
:func:`agent_team.graph.build_sqlite_checkpointer`); the call is a plain
|
||||
subprocess so this binding adds no import-time dependency on the
|
||||
orchestrator's package graph.
|
||||
def _invoke(instruction: str) -> str:
|
||||
from models import get_fast_coder # noqa: PLC0415 - intentional deferred import
|
||||
|
||||
SECURITY: this subprocess only ASKS the model for diff text — it is a model
|
||||
invocation, not a patch application. It does not run ``git``, does not apply
|
||||
anything, and does not touch the target repo. Its stdout is untrusted input
|
||||
handed back to :func:`build_candidate_diff` for defensive parsing.
|
||||
coder = get_fast_coder()
|
||||
result = coder.invoke(instruction)
|
||||
text = getattr(result, "content", result)
|
||||
return text if isinstance(text, str) else str(text)
|
||||
|
||||
return _invoke
|
||||
|
||||
|
||||
def subprocess_build(
|
||||
instruction: str, *, route: _OrchestratorRoute | None = None
|
||||
) -> str:
|
||||
"""Opt-in fallback :data:`BuildCallable`: shell out to ``run.py fast_coder``.
|
||||
|
||||
The original subprocess-based build path, retained as an opt-in alternative
|
||||
to the in-process default (:func:`make_fast_coder_invoker`). Use this in
|
||||
environments where the orchestrator's ``models`` stack is unavailable or for
|
||||
debugging.
|
||||
|
||||
SECURITY: this subprocess only ASKS the model for diff text — it does not
|
||||
run ``git``, apply anything, or touch the target repo. Its stdout is
|
||||
untrusted input handed back to :func:`build_candidate_diff` for defensive
|
||||
parsing.
|
||||
"""
|
||||
route = (
|
||||
route if route is not None else _OrchestratorRoute(root=_orchestrator_root())
|
||||
|
|
@ -188,6 +207,26 @@ def default_build(instruction: str, *, route: _OrchestratorRoute | None = None)
|
|||
return completed.stdout
|
||||
|
||||
|
||||
def default_build(instruction: str, *, route: _OrchestratorRoute | None = None) -> str:
|
||||
"""Default :data:`BuildCallable`: route the build to DeepSeek ``fast_coder`` in-process.
|
||||
|
||||
Calls DeepSeek in-process via the orchestrator's ``models.get_fast_coder()``
|
||||
(WS1). The subprocess path is available as :func:`subprocess_build` for
|
||||
environments where the models stack is absent.
|
||||
|
||||
No orchestrator module is imported at module top (deferred, mirroring
|
||||
:func:`agent_team.graph.build_sqlite_checkpointer`); any import error
|
||||
propagates to the caller (:func:`build_candidate_diff`) which wraps it as a
|
||||
``no_op`` result so the builder fails SAFE.
|
||||
|
||||
SECURITY: this call only ASKS the model for diff text — it is a model
|
||||
invocation, not a patch application. The returned text is untrusted and must
|
||||
be defensively parsed by the caller.
|
||||
"""
|
||||
invoker = make_fast_coder_invoker()
|
||||
return invoker(instruction)
|
||||
|
||||
|
||||
def _render_build_instruction(plan: Mapping[str, Any], state: Mapping[str, Any]) -> str:
|
||||
"""Render the approved plan into a mechanical-edit instruction for fast_coder.
|
||||
|
||||
|
|
|
|||
|
|
@ -221,10 +221,12 @@ def make_cross_reviewer_invoker() -> PlanReviewer:
|
|||
return _invoke
|
||||
|
||||
|
||||
# The default reviewer: subprocess to run.py (the design's "reuse run.py in
|
||||
# place" path). Built once; resolution of run.py / timeout still happens per
|
||||
# call so config and env overrides apply.
|
||||
default_plan_reviewer: PlanReviewer = make_run_py_invoker()
|
||||
# The default reviewer: in-process GPT-4.1 (WS1). The subprocess path is kept
|
||||
# as an opt-in fallback via make_run_py_invoker(). The in-process path avoids
|
||||
# the subprocess fork overhead and resolves run.py location ambiguity; the
|
||||
# subprocess path remains available for environments where the orchestrator
|
||||
# models stack is absent or for debugging.
|
||||
default_plan_reviewer: PlanReviewer = make_cross_reviewer_invoker()
|
||||
|
||||
|
||||
def review_plan(
|
||||
|
|
|
|||
|
|
@ -685,6 +685,11 @@ def _cmd_serve(args: argparse.Namespace, *, out: Any) -> int:
|
|||
job; this command owns the maintenance loop. Runs until interrupted.
|
||||
"""
|
||||
coordinator = _build_coordinator(args)
|
||||
# Bind the in-process non-Claude invokers (GPT-4.1 review, DeepSeek build)
|
||||
# alongside the Claude subscription invoker (WS1).
|
||||
from agent_team.invoker_multi import bind_multi_invoker # noqa: PLC0415
|
||||
|
||||
bind_multi_invoker()
|
||||
print("agent-team coordinator starting (Ctrl-C to stop)", file=out)
|
||||
coordinator.serve()
|
||||
return 0 # pragma: no cover - serve() loops until interrupted
|
||||
|
|
|
|||
|
|
@ -275,17 +275,18 @@ def test_source_has_no_patch_application_or_fs_write_paths() -> None:
|
|||
}, f"module uses unexpected subprocess primitives: {subprocess_attrs}"
|
||||
|
||||
|
||||
def test_default_build_invokes_run_py_as_list_argv(
|
||||
def test_subprocess_build_invokes_run_py_as_list_argv(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
"""``default_build`` shells ``run.py`` via list-form argv (no shell) + maps stdout.
|
||||
"""``subprocess_build`` (opt-in fallback) shells ``run.py`` via list-form argv (no shell).
|
||||
|
||||
Executes the subprocess path (not just AST-checks it): monkeypatches
|
||||
``subprocess.run`` to capture the invocation and return canned stdout. The
|
||||
argv MUST be the list form ``["python3", <root>/run.py, <instruction>]`` so
|
||||
the instruction can never be interpreted by a shell (no ``shell=True``), and
|
||||
the return value is the subprocess stdout verbatim.
|
||||
WS1: the subprocess path is kept as an opt-in fallback under the renamed
|
||||
``subprocess_build``. Tests it with a monkeypatched subprocess.run.
|
||||
The argv MUST be the list form ``["python3", <root>/run.py, <instruction>]``
|
||||
so the instruction can never be interpreted by a shell (no ``shell=True``).
|
||||
"""
|
||||
from agent_team.nodes.builders_llm import subprocess_build
|
||||
|
||||
captured: dict[str, Any] = {}
|
||||
|
||||
class _FakeCompleted:
|
||||
|
|
@ -300,7 +301,7 @@ def test_default_build_invokes_run_py_as_list_argv(
|
|||
|
||||
root = Path("/tmp/fake-orchestrator-root")
|
||||
route = builders_llm._OrchestratorRoute(root=root)
|
||||
out = default_build("do the edit", route=route)
|
||||
out = subprocess_build("do the edit", route=route)
|
||||
|
||||
assert out == "DIFF-FROM-SUBPROCESS"
|
||||
# List-form argv (no shell): exactly python3, the run.py path, the instruction.
|
||||
|
|
@ -328,14 +329,17 @@ def test_as_diff_builder_empty_raises_build_error_in_real_node() -> None:
|
|||
|
||||
|
||||
def test_default_build_is_the_injection_default() -> None:
|
||||
"""The default build seam is default_build (the DeepSeek/orchestrator route)."""
|
||||
"""default_build uses in-process DeepSeek (WS1); subprocess_build kept as fallback."""
|
||||
sig = inspect.signature(build_candidate_diff)
|
||||
assert sig.parameters["build"].default is None
|
||||
# default_build is what gets used when build is None — assert it's callable
|
||||
# and routes to a subprocess to run.py (string check, no execution).
|
||||
# default_build delegates to make_fast_coder_invoker (in-process, WS1);
|
||||
# subprocess_build retains the old subprocess path as an opt-in fallback.
|
||||
src = inspect.getsource(default_build)
|
||||
assert "run.py" in src
|
||||
assert "subprocess.run" in src
|
||||
assert "make_fast_coder_invoker" in src
|
||||
from agent_team.nodes.builders_llm import subprocess_build
|
||||
|
||||
subprocess_src = inspect.getsource(subprocess_build)
|
||||
assert "subprocess.run" in subprocess_src
|
||||
|
||||
# builders are DeepSeek (orchestrator fast_coder), NOT Claude: assert the
|
||||
# module never CALLS billing.claude_invoke (AST, so docstring mentions of the
|
||||
|
|
|
|||
433
agent-team/tests/test_ws1_invoker_multi_api.py
Normal file
433
agent-team/tests/test_ws1_invoker_multi_api.py
Normal file
|
|
@ -0,0 +1,433 @@
|
|||
"""WS1 tests: invoker_multi dispatch + api.py HTTP API seam.
|
||||
|
||||
Covers:
|
||||
* invoker_multi.multi_invoke dispatches to the correct model factory (mocked)
|
||||
* invoker_multi.make_fast_coder_invoker / make_scanner_invoker return callables
|
||||
* invoker_multi.bind_multi_invoker wires the review loop seam
|
||||
* api.py: 401 without bearer token; 200 with valid token (endpoints stubbed)
|
||||
* api.py: make_app returns a FastAPI app; GET /tasks 404 when no state
|
||||
* review_loop_llm default_plan_reviewer is now make_cross_reviewer_invoker (in-process)
|
||||
* builders_llm subprocess_build is kept as opt-in; make_fast_coder_invoker callable
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import importlib.util
|
||||
import sys
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
import pytest
|
||||
|
||||
_AGENT_TEAM_DIR = Path(__file__).resolve().parents[1]
|
||||
if str(_AGENT_TEAM_DIR) not in sys.path:
|
||||
sys.path.insert(0, str(_AGENT_TEAM_DIR))
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# invoker_multi: basic dispatch via mocked models
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def _make_fake_model(return_text: str = "model output") -> MagicMock:
|
||||
"""Return a mock model whose .invoke() yields a fake response."""
|
||||
mock = MagicMock()
|
||||
result = MagicMock()
|
||||
result.content = return_text
|
||||
mock.invoke.return_value = result
|
||||
return mock
|
||||
|
||||
|
||||
def test_multi_invoke_cross_reviewer(monkeypatch: Any) -> None:
|
||||
from agent_team import invoker_multi
|
||||
|
||||
fake_model = _make_fake_model("APPROVE verdict")
|
||||
fake_models = MagicMock()
|
||||
fake_models.get_cross_reviewer.return_value = fake_model
|
||||
fake_models.get_fast_coder.return_value = MagicMock()
|
||||
fake_models.get_scanner.return_value = MagicMock()
|
||||
|
||||
with patch.dict(sys.modules, {"models": fake_models}):
|
||||
result = invoker_multi.multi_invoke("review this plan", model="cross_reviewer")
|
||||
|
||||
assert result == "APPROVE verdict"
|
||||
fake_model.invoke.assert_called_once_with("review this plan")
|
||||
|
||||
|
||||
def test_multi_invoke_fast_coder(monkeypatch: Any) -> None:
|
||||
from agent_team import invoker_multi
|
||||
|
||||
fake_model = _make_fake_model("--- a/file.py\n+++ b/file.py\n@@")
|
||||
fake_models = MagicMock()
|
||||
fake_models.get_fast_coder.return_value = fake_model
|
||||
fake_models.get_cross_reviewer.return_value = MagicMock()
|
||||
fake_models.get_scanner.return_value = MagicMock()
|
||||
|
||||
with patch.dict(sys.modules, {"models": fake_models}):
|
||||
result = invoker_multi.multi_invoke("add a login form", model="fast_coder")
|
||||
|
||||
assert "@@" in result
|
||||
|
||||
|
||||
def test_multi_invoke_scanner(monkeypatch: Any) -> None:
|
||||
from agent_team import invoker_multi
|
||||
|
||||
fake_model = _make_fake_model("no issues found")
|
||||
fake_models = MagicMock()
|
||||
fake_models.get_scanner.return_value = fake_model
|
||||
fake_models.get_cross_reviewer.return_value = MagicMock()
|
||||
fake_models.get_fast_coder.return_value = MagicMock()
|
||||
|
||||
with patch.dict(sys.modules, {"models": fake_models}):
|
||||
result = invoker_multi.multi_invoke("scan the diff", model="scanner")
|
||||
|
||||
assert result == "no issues found"
|
||||
|
||||
|
||||
def test_multi_invoke_unknown_model() -> None:
|
||||
from agent_team import invoker_multi
|
||||
|
||||
fake_models = MagicMock()
|
||||
with patch.dict(sys.modules, {"models": fake_models}):
|
||||
with pytest.raises(ValueError, match="unknown model key"):
|
||||
invoker_multi.multi_invoke("hi", model="gpt_4_turbo")
|
||||
|
||||
|
||||
def test_multi_invoke_models_unavailable() -> None:
|
||||
from agent_team import invoker_multi
|
||||
|
||||
saved = sys.modules.pop("models", None)
|
||||
try:
|
||||
# Ensure "models" is NOT in sys.modules so the import inside multi_invoke fails.
|
||||
# We also need to make sure it's not importable from the path.
|
||||
with pytest.raises(RuntimeError, match="Cannot import orchestrator"):
|
||||
# patch.dict with a None value keeps it out of sys.modules
|
||||
with patch.dict(sys.modules, {"models": None}): # type: ignore[dict-item]
|
||||
invoker_multi.multi_invoke("hi", model="cross_reviewer")
|
||||
finally:
|
||||
if saved is not None:
|
||||
sys.modules["models"] = saved
|
||||
|
||||
|
||||
def test_make_fast_coder_invoker_returns_callable() -> None:
|
||||
from agent_team.invoker_multi import make_fast_coder_invoker
|
||||
|
||||
invoker = make_fast_coder_invoker()
|
||||
assert callable(invoker)
|
||||
|
||||
|
||||
def test_make_scanner_invoker_returns_callable() -> None:
|
||||
from agent_team.invoker_multi import make_scanner_invoker
|
||||
|
||||
invoker = make_scanner_invoker()
|
||||
assert callable(invoker)
|
||||
|
||||
|
||||
def test_make_fast_coder_invoker_calls_multi_invoke() -> None:
|
||||
from agent_team import invoker_multi
|
||||
|
||||
fake_model = _make_fake_model("the diff")
|
||||
fake_models = MagicMock()
|
||||
fake_models.get_fast_coder.return_value = fake_model
|
||||
fake_models.get_cross_reviewer.return_value = MagicMock()
|
||||
fake_models.get_scanner.return_value = MagicMock()
|
||||
|
||||
invoker = invoker_multi.make_fast_coder_invoker()
|
||||
with patch.dict(sys.modules, {"models": fake_models}):
|
||||
result = invoker("build instruction")
|
||||
|
||||
assert result == "the diff"
|
||||
|
||||
|
||||
def test_bind_multi_invoker_sets_review_seam() -> None:
|
||||
from agent_team import invoker_multi
|
||||
from agent_team.nodes import review_loop
|
||||
|
||||
original = review_loop._review_invoker # save
|
||||
try:
|
||||
fake_model = _make_fake_model("APPROVE")
|
||||
fake_models = MagicMock()
|
||||
fake_models.get_cross_reviewer.return_value = fake_model
|
||||
fake_models.get_fast_coder.return_value = MagicMock()
|
||||
fake_models.get_scanner.return_value = MagicMock()
|
||||
|
||||
with patch.dict(sys.modules, {"models": fake_models}):
|
||||
invoker_multi.bind_multi_invoker()
|
||||
|
||||
# The review loop seam should now be the cross-reviewer invoker.
|
||||
assert review_loop._review_invoker is not original
|
||||
finally:
|
||||
review_loop._review_invoker = original # restore
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# review_loop_llm: default is now in-process, subprocess kept as opt-in
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def test_review_loop_llm_default_is_cross_reviewer_invoker() -> None:
|
||||
from agent_team.nodes import review_loop_llm
|
||||
|
||||
# The module-level default should be the closure from make_cross_reviewer_invoker,
|
||||
# NOT the subprocess-based make_run_py_invoker closure. We test by name.
|
||||
# Both are closures, so we check the closure cell names via __code__.
|
||||
fn = review_loop_llm.default_plan_reviewer
|
||||
assert callable(fn)
|
||||
# make_cross_reviewer_invoker creates a _invoke that uses get_cross_reviewer
|
||||
# internally; make_run_py_invoker creates one that uses subprocess.run.
|
||||
# We verify by checking that the default raises RuntimeError (models absent)
|
||||
# NOT FileNotFoundError (run.py absent), confirming it's the in-process path.
|
||||
saved = sys.modules.pop("models", None)
|
||||
try:
|
||||
with patch.dict(sys.modules, {"models": None}): # type: ignore[dict-item]
|
||||
with pytest.raises(
|
||||
(RuntimeError, AttributeError, TypeError, ModuleNotFoundError)
|
||||
):
|
||||
fn("any prompt")
|
||||
finally:
|
||||
if saved is not None:
|
||||
sys.modules["models"] = saved
|
||||
|
||||
|
||||
def test_review_loop_llm_make_run_py_invoker_still_available() -> None:
|
||||
from agent_team.nodes.review_loop_llm import make_run_py_invoker
|
||||
|
||||
fn = make_run_py_invoker()
|
||||
assert callable(fn)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# builders_llm: subprocess_build opt-in; make_fast_coder_invoker available
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def test_builders_llm_subprocess_build_available() -> None:
|
||||
from agent_team.nodes.builders_llm import subprocess_build
|
||||
|
||||
assert callable(subprocess_build)
|
||||
|
||||
|
||||
def test_builders_llm_make_fast_coder_invoker_available() -> None:
|
||||
from agent_team.nodes.builders_llm import make_fast_coder_invoker
|
||||
|
||||
fn = make_fast_coder_invoker()
|
||||
assert callable(fn)
|
||||
|
||||
|
||||
def test_builders_llm_default_build_uses_in_process() -> None:
|
||||
from agent_team.nodes import builders_llm
|
||||
|
||||
fake_model = _make_fake_model("--- a/x.py\n+++ b/x.py\n@@ -1 +1 @@\n-old\n+new")
|
||||
fake_models = MagicMock()
|
||||
fake_models.get_fast_coder.return_value = fake_model
|
||||
fake_models.get_cross_reviewer.return_value = MagicMock()
|
||||
fake_models.get_scanner.return_value = MagicMock()
|
||||
|
||||
with patch.dict(sys.modules, {"models": fake_models}):
|
||||
result = builders_llm.default_build("do the thing")
|
||||
|
||||
assert "@@" in result
|
||||
fake_model.invoke.assert_called_once()
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# api.py: bearer-token auth + endpoint smoke tests via FastAPI TestClient
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
# fastapi is a box-only dependency (installed on the R720 during attended
|
||||
# deploy, never in CI or the root requirements). api.py imports it lazily, so
|
||||
# these TestClient-based smoke tests are the only thing that hard-requires it.
|
||||
# Skip them when it is absent rather than failing collection in CI.
|
||||
requires_fastapi = pytest.mark.skipif(
|
||||
importlib.util.find_spec("fastapi") is None,
|
||||
reason="fastapi not installed (box-only dependency)",
|
||||
)
|
||||
|
||||
|
||||
@pytest.fixture()
|
||||
def api_token(monkeypatch: Any) -> str:
|
||||
token = "test-secret-token-for-ws1"
|
||||
monkeypatch.setenv("AGENT_TEAM_API_TOKEN", token)
|
||||
return token
|
||||
|
||||
|
||||
def _make_stub_coordinator(thread_id: str = "t-test-123") -> MagicMock:
|
||||
"""Return a coordinator stub that start_task returns a fixed thread_id."""
|
||||
coord = MagicMock()
|
||||
coord.start_task.return_value = thread_id
|
||||
coord.graph = None # no graph for simple smoke tests
|
||||
coord.setup = MagicMock()
|
||||
return coord
|
||||
|
||||
|
||||
@requires_fastapi
|
||||
def test_api_401_without_token(api_token: str) -> None:
|
||||
from fastapi.testclient import TestClient
|
||||
|
||||
from agent_team.api import make_app
|
||||
|
||||
app = make_app(coordinator=_make_stub_coordinator())
|
||||
client = TestClient(app, raise_server_exceptions=False)
|
||||
resp = client.post("/tasks", json={"task": "add login"})
|
||||
assert resp.status_code == 401
|
||||
|
||||
|
||||
@requires_fastapi
|
||||
def test_api_401_wrong_token(api_token: str) -> None:
|
||||
from fastapi.testclient import TestClient
|
||||
|
||||
from agent_team.api import make_app
|
||||
|
||||
app = make_app(coordinator=_make_stub_coordinator())
|
||||
client = TestClient(app, raise_server_exceptions=False)
|
||||
resp = client.post(
|
||||
"/tasks",
|
||||
json={"task": "add login"},
|
||||
headers={"Authorization": "Bearer wrong-token"},
|
||||
)
|
||||
assert resp.status_code == 401
|
||||
|
||||
|
||||
@requires_fastapi
|
||||
def test_api_post_tasks_200(api_token: str) -> None:
|
||||
from fastapi.testclient import TestClient
|
||||
|
||||
from agent_team.api import make_app
|
||||
|
||||
stub = _make_stub_coordinator("thread-abc")
|
||||
app = make_app(coordinator=stub)
|
||||
client = TestClient(app)
|
||||
resp = client.post(
|
||||
"/tasks",
|
||||
json={"task": "add a login form"},
|
||||
headers={"Authorization": f"Bearer {api_token}"},
|
||||
)
|
||||
assert resp.status_code == 200, resp.text
|
||||
assert resp.json()["thread_id"] == "thread-abc"
|
||||
stub.start_task.assert_called_once()
|
||||
|
||||
|
||||
@requires_fastapi
|
||||
def test_api_get_task_404_no_state(api_token: str) -> None:
|
||||
from fastapi.testclient import TestClient
|
||||
|
||||
from agent_team.api import make_app
|
||||
|
||||
stub = _make_stub_coordinator()
|
||||
app = make_app(coordinator=stub)
|
||||
client = TestClient(app, raise_server_exceptions=False)
|
||||
# graph is None -> 503 not ready
|
||||
resp = client.get(
|
||||
"/tasks/t-missing",
|
||||
headers={"Authorization": f"Bearer {api_token}"},
|
||||
)
|
||||
assert resp.status_code == 503
|
||||
|
||||
|
||||
@requires_fastapi
|
||||
def test_api_get_task_with_graph_state(api_token: str) -> None:
|
||||
from fastapi.testclient import TestClient
|
||||
|
||||
from agent_team.api import make_app
|
||||
|
||||
stub = _make_stub_coordinator()
|
||||
# Attach a mock graph that returns a snapshot with values
|
||||
snapshot = MagicMock()
|
||||
snapshot.values = {"task": "add login", "current_phase": "CLARIFY"}
|
||||
stub.graph = MagicMock()
|
||||
stub.graph.get_state.return_value = snapshot
|
||||
app = make_app(coordinator=stub)
|
||||
client = TestClient(app)
|
||||
resp = client.get(
|
||||
"/tasks/t-abc",
|
||||
headers={"Authorization": f"Bearer {api_token}"},
|
||||
)
|
||||
assert resp.status_code == 200
|
||||
data = resp.json()
|
||||
assert data["thread_id"] == "t-abc"
|
||||
assert data["state"]["task"] == "add login"
|
||||
|
||||
|
||||
@requires_fastapi
|
||||
def test_api_orchestrator_invoke_calls_subprocess(
|
||||
api_token: str, tmp_path: Any
|
||||
) -> None:
|
||||
from fastapi.testclient import TestClient
|
||||
|
||||
from agent_team.api import make_app
|
||||
|
||||
# Write a fake run.py that echos the first argv
|
||||
fake_run = tmp_path / "run.py"
|
||||
fake_run.write_text("import sys; print(sys.argv[1])")
|
||||
|
||||
stub = _make_stub_coordinator()
|
||||
app = make_app(coordinator=stub)
|
||||
|
||||
# Patch _resolve_run_py to return our fake script
|
||||
with patch("agent_team.api._resolve_run_py", return_value=fake_run):
|
||||
client = TestClient(app)
|
||||
resp = client.post(
|
||||
"/orchestrator/invoke",
|
||||
json={"prompt": "hello world"},
|
||||
headers={"Authorization": f"Bearer {api_token}"},
|
||||
)
|
||||
|
||||
assert resp.status_code == 200
|
||||
assert "hello world" in resp.json()["text"]
|
||||
|
||||
|
||||
@requires_fastapi
|
||||
def test_api_no_token_env_raises_at_build(monkeypatch: Any) -> None:
|
||||
from agent_team.api import make_app
|
||||
|
||||
monkeypatch.delenv("AGENT_TEAM_API_TOKEN", raising=False)
|
||||
# Fail-fast: an unconfigured token must raise at app-build time, not serve a
|
||||
# request first (the eager _get_token() in make_app).
|
||||
with pytest.raises(RuntimeError, match="AGENT_TEAM_API_TOKEN"):
|
||||
make_app(coordinator=_make_stub_coordinator())
|
||||
|
||||
|
||||
@requires_fastapi
|
||||
def test_api_docs_and_openapi_disabled(api_token: str) -> None:
|
||||
from fastapi.testclient import TestClient
|
||||
|
||||
from agent_team.api import make_app
|
||||
|
||||
app = make_app(coordinator=_make_stub_coordinator())
|
||||
client = TestClient(app, raise_server_exceptions=False)
|
||||
# The unauthenticated docs/schema routes must be disabled (VPN-only API).
|
||||
for path in ("/docs", "/redoc", "/openapi.json"):
|
||||
assert client.get(path).status_code == 404, path
|
||||
|
||||
|
||||
@requires_fastapi
|
||||
def test_api_invoke_concurrency_cap_returns_429(api_token: str, tmp_path: Any) -> None:
|
||||
import agent_team.api as api_mod
|
||||
from fastapi.testclient import TestClient
|
||||
|
||||
from agent_team.api import make_app
|
||||
|
||||
fake_run = tmp_path / "run.py"
|
||||
fake_run.write_text("import sys; print(sys.argv[1])")
|
||||
app = make_app(coordinator=_make_stub_coordinator())
|
||||
|
||||
# Exhaust the bounded semaphore so the request sees no free slot → 429.
|
||||
acquired = [
|
||||
api_mod._invoke_semaphore.acquire(blocking=False)
|
||||
for _ in range(api_mod._MAX_CONCURRENT_INVOKES)
|
||||
]
|
||||
try:
|
||||
assert all(acquired)
|
||||
with patch("agent_team.api._resolve_run_py", return_value=fake_run):
|
||||
client = TestClient(app, raise_server_exceptions=False)
|
||||
resp = client.post(
|
||||
"/orchestrator/invoke",
|
||||
json={"prompt": "hello"},
|
||||
headers={"Authorization": f"Bearer {api_token}"},
|
||||
)
|
||||
assert resp.status_code == 429
|
||||
finally:
|
||||
for ok in acquired:
|
||||
if ok:
|
||||
api_mod._invoke_semaphore.release()
|
||||
|
|
@ -7,3 +7,6 @@ langchain-google-genai==4.2.5
|
|||
langchain-community==0.4.2
|
||||
composio-langgraph==0.15.0
|
||||
python-dotenv==1.2.2
|
||||
# WS1 agent-team HTTP API (agent_team/api.py): FastAPI app + uvicorn ASGI server.
|
||||
fastapi==0.136.1
|
||||
uvicorn==0.46.0
|
||||
|
|
|
|||
Reference in a new issue