diff --git a/agent-team/agent_team/api.py b/agent-team/agent_team/api.py new file mode 100644 index 0000000..8ba955c --- /dev/null +++ b/agent-team/agent_team/api.py @@ -0,0 +1,293 @@ +"""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 `` + 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 ``/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) diff --git a/agent-team/agent_team/invoker_multi.py b/agent-team/agent_team/invoker_multi.py new file mode 100644 index 0000000..a34ffd8 --- /dev/null +++ b/agent-team/agent_team/invoker_multi.py @@ -0,0 +1,171 @@ +"""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 ``/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()) diff --git a/agent-team/agent_team/nodes/builders_llm.py b/agent-team/agent_team/nodes/builders_llm.py index 30701d6..5440e11 100644 --- a/agent-team/agent_team/nodes/builders_llm.py +++ b/agent-team/agent_team/nodes/builders_llm.py @@ -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,39 @@ 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 /run.py - ""`` — 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 +205,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. diff --git a/agent-team/agent_team/nodes/review_loop_llm.py b/agent-team/agent_team/nodes/review_loop_llm.py index 1a9a8ed..3d8de87 100644 --- a/agent-team/agent_team/nodes/review_loop_llm.py +++ b/agent-team/agent_team/nodes/review_loop_llm.py @@ -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( diff --git a/agent-team/run-team.py b/agent-team/run-team.py index 2986149..7d07e3a 100644 --- a/agent-team/run-team.py +++ b/agent-team/run-team.py @@ -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 diff --git a/agent-team/tests/test_builders_llm.py b/agent-team/tests/test_builders_llm.py index 55bb366..cacfc46 100644 --- a/agent-team/tests/test_builders_llm.py +++ b/agent-team/tests/test_builders_llm.py @@ -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", /run.py, ]`` 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", /run.py, ]`` + 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,16 @@ 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 diff --git a/agent-team/tests/test_ws1_invoker_multi_api.py b/agent-team/tests/test_ws1_invoker_multi_api.py new file mode 100644 index 0000000..f6c7786 --- /dev/null +++ b/agent-team/tests/test_ws1_invoker_multi_api.py @@ -0,0 +1,375 @@ +"""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 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 +# --------------------------------------------------------------------------- + + +@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 + + +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 + + +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 + + +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() + + +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 + + +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" + + +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"] + + +def test_api_no_token_env_raises_on_request(monkeypatch: Any) -> None: + from fastapi.testclient import TestClient + + from agent_team.api import make_app + + monkeypatch.delenv("AGENT_TEAM_API_TOKEN", raising=False) + app = make_app(coordinator=_make_stub_coordinator()) + client = TestClient(app, raise_server_exceptions=False) + # A missing env var raises RuntimeError inside the dependency — should 500 + # or 401; either way the request must not succeed. + resp = client.post( + "/tasks", + json={"task": "hi"}, + headers={"Authorization": "Bearer anything"}, + ) + assert resp.status_code in (401, 500)