From 41132a46155ca85de553ab69a3e3fa14f5e13c08 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Thu, 18 Jun 2026 15:30:07 -0400 Subject: [PATCH] feat(agent-team): read-only CI-result fetcher for P3 verify gate (opt-in, inert) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ci_fetcher.py: fail-closed CiResultFetcher reading the GitHub Actions run conclusion via a read-only PAT (AGENT_TEAM_CI_READ_TOKEN→GITHUB_TOKEN), returns {run_id,conclusion,diff_hash} or None on any error. Data-fetcher only — ci_gate owns the verdict; never writes, no OIDC/AWS, never reads patch artifacts. coordinator gains opt-in gated_build_verify_wiring() composing it via bind_ci_result_fetcher; NOT wired into the default run-team.py path. 20 tests. --- agent-team/agent_team/ci_fetcher.py | 260 +++++++++++++++++++++++++++ agent-team/agent_team/coordinator.py | 57 ++++++ agent-team/tests/test_ci_fetcher.py | 202 +++++++++++++++++++++ 3 files changed, 519 insertions(+) create mode 100644 agent-team/agent_team/ci_fetcher.py create mode 100644 agent-team/tests/test_ci_fetcher.py diff --git a/agent-team/agent_team/ci_fetcher.py b/agent-team/agent_team/ci_fetcher.py new file mode 100644 index 0000000..6ac298c --- /dev/null +++ b/agent-team/agent_team/ci_fetcher.py @@ -0,0 +1,260 @@ +"""Real, read-only CI-result fetcher — the GATED-LIVE seam for VERIFY (§3.3.2 #4). + +This module implements the production :data:`~agent_team.nodes.build_verify_subgraph.CiResultFetcher` +that the build->verify subgraph's VERIFY node injects once the §3.3.2 CI +trust-boundary clears ``/sh-security-review`` + the GPT-4.1 cross-review. Until +then the subgraph runs with the INERT default fetcher (``_no_ci_result`` -> the +gate BLOCKs and the task parks); binding THIS fetcher only gives the pure-code +gate (:func:`agent_team.ci_gate.evaluate_ci_gate`) an authenticated conclusion +to read — it never makes the LLM the pass authority. + +What it is (and, just as importantly, what it is NOT): + +* It is a **pure DATA fetcher.** Given the task ``state``, it reads + ``state["run_id"]`` (and echoes ``state.get("diff_hash")``), calls the GitHub + Actions REST API ``GET /repos/{owner}/{repo}/actions/runs/{run_id}`` with a + **READ-ONLY** token, and returns the authenticated run ``conclusion`` as a + mapping ``{"run_id", "conclusion", "diff_hash"}``. The gate owns the verdict; + this module never derives pass/fail itself. +* It is **fail-closed.** ANY error — missing ``run_id``, missing token, 404, + auth failure, malformed JSON, a network/timeout error, an unexpected status — + returns ``None``. A ``None`` result makes the gate BLOCK (never a silent + pass), so a broken fetch parks the task for a human rather than shipping. +* It **never writes.** No ``POST``/``PATCH``, no ``git``, no patch apply, no + filesystem mutation. It performs exactly one read-only GET. +* It uses **no cloud credentials and no token-federation** — only a GitHub + read-only token, resolved at CALL time from the environment (matching the + prevailing transport idiom in :mod:`agent_team.transport.github_live`). It + NEVER reads a success/failure file the patch could have written (boundary 4). + +Prevailing HTTP approach: mirrors :mod:`agent_team.transport.github_live` — a +thin ``requests`` session, deferred import (``requests`` is optional and may be +absent pre-deploy), call-time token resolution, and an injectable ``client`` for +testability. Unlike the poster, a missing token / missing ``requests`` here does +NOT raise: it fails closed to ``None`` so the gate BLOCKs (the fetcher's whole +contract is "no authenticated result -> None"). The dedicated read-only env var +``AGENT_TEAM_CI_READ_TOKEN`` is preferred, falling back to ``GITHUB_TOKEN``. +""" + +from __future__ import annotations + +import logging +import os +from collections.abc import Mapping +from typing import Any + +from agent_team.nodes.build_verify_subgraph import CiResultFetcher + +__all__ = [ + "CI_READ_TOKEN_ENV", + "GITHUB_API_ROOT", + "build_ci_result_fetcher", + "fetch_ci_result", +] + +_LOG = logging.getLogger("agent_team.ci_fetcher") + +# The dedicated read-only token env var (preferred), falling back to the generic +# GITHUB_TOKEN (the idiom github_live uses). The token MUST be read-only — the +# fetcher only ever GETs; a write-scoped token here would be unnecessary blast +# radius (provisioning issues a read-only fine-grained PAT / read-only var). +CI_READ_TOKEN_ENV = "AGENT_TEAM_CI_READ_TOKEN" +_FALLBACK_TOKEN_ENV = "GITHUB_TOKEN" + +GITHUB_API_ROOT = "https://api.github.com" + +# Conservative default timeout for the single read-only GET. A hang must fail +# closed (-> None -> gate BLOCK), never wedge the verifier. +_DEFAULT_TIMEOUT_S = 15.0 + + +def _resolve_read_token() -> str | None: + """Resolve the read-only GitHub token at call time, or ``None``. + + Prefers :data:`CI_READ_TOKEN_ENV`, falls back to ``GITHUB_TOKEN``. Returns + ``None`` when neither is set so the fetcher fails closed (the caller maps a + missing token to a ``None`` result -> gate BLOCK), rather than raising. + """ + return os.environ.get(CI_READ_TOKEN_ENV) or os.environ.get(_FALLBACK_TOKEN_ENV) + + +def _build_session(token: str) -> Any: + """Lazily construct a read-only ``requests.Session`` (deferred optional import). + + Mirrors :func:`agent_team.transport.github_live._build_session` (deferred + ``requests`` import, ``Authorization: Bearer`` + the API-version header). The + session is used for exactly one GET; no write verbs are ever issued. + + Raises :class:`RuntimeError` only if ``requests`` is unavailable — the caller + catches it and fails closed to ``None`` so a pre-deploy environment without + the optional dependency simply BLOCKs (never a spurious pass). + """ + import requests # deferred: optional dependency (see module docstring) + + session = requests.Session() + session.headers.update( + { + "Authorization": f"Bearer {token}", + "Accept": "application/vnd.github+json", + "X-GitHub-Api-Version": "2022-11-28", + } + ) + return session + + +def fetch_ci_result( + state: Mapping[str, Any], + *, + owner: str, + repo: str, + client: Any = None, + api_root: str = GITHUB_API_ROOT, + timeout: float = _DEFAULT_TIMEOUT_S, +) -> dict[str, Any] | None: + """Fetch the authenticated CI run conclusion as DATA, or ``None`` (fail-closed). + + Reads ``state["run_id"]`` (required) and echoes ``state.get("diff_hash")``, + then GETs ``{api_root}/repos/{owner}/{repo}/actions/runs/{run_id}`` with the + read-only token and returns:: + + {"run_id": , "conclusion": , "diff_hash": } + + Returns ``None`` on ANY failure — no ``run_id`` in state, no resolvable + token, ``requests`` unavailable, a non-2xx status (404/401/403/...), a + malformed/absent JSON body, a missing ``conclusion``, or a network/timeout + error. ``None`` is the fail-closed signal: the pure-code gate treats it as + "no authenticated result" and BLOCKs, so a broken fetch parks the task. This + function NEVER derives the verdict (the gate owns pass/fail) and NEVER + writes. + + ``client`` injects a pre-built ``requests``-like session for tests (any + object with ``get(url, *, timeout) -> response`` exposing ``status_code`` + and ``json()``). When omitted, a read-only session is built from the + resolved token. + """ + run_id = state.get("run_id") + if run_id is None or str(run_id) == "": + _LOG.warning("ci_fetcher: no run_id in state; failing closed to None") + return None + run_id = str(run_id) + + echoed_hash = state.get("diff_hash") + + if client is None: + token = _resolve_read_token() + if not token: + _LOG.warning( + "ci_fetcher: no read-only token (%s/%s) set; failing closed to None", + CI_READ_TOKEN_ENV, + _FALLBACK_TOKEN_ENV, + ) + return None + try: + client = _build_session(token) + except Exception: # noqa: BLE001 - any build failure fails closed + _LOG.warning("ci_fetcher: could not build HTTP client; failing closed") + return None + + url = f"{api_root.rstrip('/')}/repos/{owner}/{repo}/actions/runs/{run_id}" + + try: + response = client.get(url, timeout=timeout) + except Exception: # noqa: BLE001 - timeout / connection / any -> fail closed + _LOG.warning("ci_fetcher: GET %s failed (network/timeout); failing closed", url) + return None + + status = _status_of(response) + if status is None or not (200 <= status < 300): + _LOG.warning("ci_fetcher: run GET returned status %r; failing closed", status) + return None + + data = _json_of(response) + if not isinstance(data, dict): + _LOG.warning("ci_fetcher: run body was not a JSON object; failing closed") + return None + + conclusion = data.get("conclusion") + # An in-progress run has conclusion=None; that is NOT an authenticated + # verdict, so fail closed (the gate would BLOCK on it anyway, but returning + # None keeps the "no result" contract clean and avoids echoing a non-verdict). + if conclusion is None or not isinstance(conclusion, str): + _LOG.info("ci_fetcher: run %s has no conclusion yet; failing closed", run_id) + return None + + # Bind the returned run_id to the run actually fetched (the API echoes id); + # fall back to the requested run_id. The gate independently re-checks this + # against expected_run_id, so this is provenance, not the trust decision. + fetched_id = data.get("id") + result_run_id = str(fetched_id) if fetched_id is not None else run_id + + return { + "run_id": result_run_id, + "conclusion": conclusion, + "diff_hash": echoed_hash, + } + + +def build_ci_result_fetcher( + *, + owner: str, + repo: str, + client: Any = None, + api_root: str = GITHUB_API_ROOT, + timeout: float = _DEFAULT_TIMEOUT_S, +) -> CiResultFetcher: + """Build a :data:`CiResultFetcher` bound to ``owner``/``repo`` (GATED-LIVE seam). + + Returns a single-argument ``state -> mapping | None`` callable shaped exactly + like the VERIFY node's injected ``ci_result_fetcher`` seam, closing over the + target ``owner``/``repo`` (and the optional injected ``client`` / ``api_root`` + / ``timeout``). Compose it with + :func:`agent_team.nodes.build_verify_subgraph.bind_ci_result_fetcher` (or the + opt-in :func:`agent_team.coordinator.gated_build_verify_wiring`) once the + §3.3.2 gate clears. It is read-only and fails closed (see + :func:`fetch_ci_result`); binding it does not enable any apply/verify + behaviour, it only gives the gate an authenticated conclusion to read. + """ + + def fetcher(state: Mapping[str, Any]) -> dict[str, Any] | None: + return fetch_ci_result( + state, + owner=owner, + repo=repo, + client=client, + api_root=api_root, + timeout=timeout, + ) + + return fetcher + + +def _status_of(response: Any) -> int | None: + """Read the HTTP status from a ``requests``-like response, or ``None``. + + Accepts ``status_code`` (``requests``) or ``status`` (a minimal fake). + Returns ``None`` if neither is present so the caller fails closed rather than + raising on an exotic object. + """ + for attr in ("status_code", "status"): + value = getattr(response, attr, None) + if value is not None: + try: + return int(value) + except (TypeError, ValueError): + return None + return None + + +def _json_of(response: Any) -> Any: + """Parse the JSON body of a ``requests``-like response, or ``None`` on failure. + + A malformed/absent body must fail closed (-> ``None`` -> the caller returns + ``None`` -> gate BLOCK), never raise into the verifier node. + """ + parser = getattr(response, "json", None) + if not callable(parser): + return None + try: + return parser() + except Exception: # noqa: BLE001 - any JSON decode error fails closed + return None diff --git a/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index bfe0760..eee73b0 100644 --- a/agent-team/agent_team/coordinator.py +++ b/agent-team/agent_team/coordinator.py @@ -68,6 +68,7 @@ __all__ = [ "Coordinator", "build_verify_wiring", "default_clarify_node_factory", + "gated_build_verify_wiring", ] _LOG = logging.getLogger("agent_team.coordinator") @@ -254,6 +255,62 @@ def build_verify_wiring() -> tuple[ return build_node, verify_node, route_after_verify +def gated_build_verify_wiring( + *, + owner: str, + repo: str, + expected_run_id: str, + allowed_scope: list[str] | None = None, + diff_builder: Any = None, + ci_client: Any = None, +) -> tuple[ + Callable[[PipelineState], "dict[str, Any]"], + Callable[[PipelineState], PipelineState], + Callable[[PipelineState], str], +]: + """Compose the GATED-LIVE P3 build->verify subgraph (OPT-IN; held for the gate). + + The live counterpart to :func:`build_verify_wiring`: it binds the REAL + read-only CI-result fetcher (:func:`agent_team.ci_fetcher.build_ci_result_fetcher`) + onto the VERIFY node so the pure-code gate has an authenticated conclusion to + read, and (optionally) a real ``diff_builder`` onto the BUILD node. The gate + still owns pass/fail — binding the fetcher only GIVES it a result, it can + never make the LLM the pass authority. + + **This is OPT-IN and is deliberately NOT wired into the default + ``run-team.py`` / ``serve`` path** (exactly like :func:`build_verify_wiring`). + It is the seam a leaf calls AFTER the §3.3.2 CI trust boundary clears + ``/sh-security-review`` + the GPT-4.1 cross-review and the GitHub App + + ``agent-apply`` environment are provisioned. Until then, importing or holding + this function changes nothing: the production coordinator passes + ``build_verify_wiring=None``, so no P3 subgraph is assembled at all. + + The CI fetcher is read-only and fails closed: a missing token / 404 / auth + failure / malformed body / timeout yields ``None`` and the gate BLOCKs (the + task parks). ``ci_client`` injects a test double; the real path builds a + read-only ``requests`` session at call time from the read-only token env var. + + Lazy-imported (ci_fetcher pulls the subgraph + verifier leaves) for the same + import-hygiene reason as the other factories. + """ + from agent_team.ci_fetcher import build_ci_result_fetcher + from agent_team.nodes.build_verify_subgraph import ( + make_build_node, + make_verify_node, + route_after_verify, + ) + from agent_team.nodes.verifier import VerifierConfig + + config = VerifierConfig( + expected_run_id=expected_run_id, + allowed_scope=allowed_scope, + ) + fetcher = build_ci_result_fetcher(owner=owner, repo=repo, client=ci_client) + build_node = make_build_node(diff_builder=diff_builder) + verify_node = make_verify_node(config, ci_result_fetcher=fetcher) + return build_node, verify_node, route_after_verify + + class Coordinator: """Owns the live Plane-2 runtime: graph + resume worker + transport (§3.3). diff --git a/agent-team/tests/test_ci_fetcher.py b/agent-team/tests/test_ci_fetcher.py new file mode 100644 index 0000000..f7f35d7 --- /dev/null +++ b/agent-team/tests/test_ci_fetcher.py @@ -0,0 +1,202 @@ +"""Unit tests for agent_team.ci_fetcher — the read-only, fail-closed CI fetcher. + +The fetcher is a pure DATA seam: on a clean authenticated run it returns the +``{run_id, conclusion, diff_hash}`` mapping the pure-code gate consumes; on ANY +error (404 / auth fail / malformed body / timeout / missing run_id / missing +token) it returns ``None`` so the gate BLOCKs (fail-closed). It NEVER derives a +verdict and NEVER writes. Everything is mocked with a fake HTTP client; no +network, no real token. +""" + +from __future__ import annotations + +import pytest + +from agent_team.ci_fetcher import ( + CI_READ_TOKEN_ENV, + build_ci_result_fetcher, + fetch_ci_result, +) + + +class _FakeResponse: + def __init__(self, status: int, body): + self.status_code = status + self._body = body + + def json(self): + if isinstance(self._body, Exception): + raise self._body + return self._body + + +class _FakeClient: + """A requests-like client recording calls; raises only write verbs if asked.""" + + def __init__(self, response=None, raise_exc=None): + self._response = response + self._raise = raise_exc + self.calls: list[tuple[str, dict]] = [] + + def get(self, url, *, timeout=None): + self.calls.append((url, {"timeout": timeout})) + if self._raise is not None: + raise self._raise + return self._response + + # If the fetcher ever tried to write, these would record it — they must not + # be called (the fetcher is read-only). + def post(self, *a, **k): # pragma: no cover - must never be called + raise AssertionError("ci_fetcher must never POST") + + def patch(self, *a, **k): # pragma: no cover - must never be called + raise AssertionError("ci_fetcher must never PATCH") + + +def _state(run_id="123", diff_hash="abc123"): + return {"run_id": run_id, "diff_hash": diff_hash} + + +# --------------------------------------------------------------------------- # +# Success path +# --------------------------------------------------------------------------- # + + +def test_success_returns_correct_mapping() -> None: + client = _FakeClient( + _FakeResponse(200, {"id": 123, "conclusion": "success", "status": "completed"}) + ) + out = fetch_ci_result(_state(), owner="o", repo="r", client=client) + assert out == {"run_id": "123", "conclusion": "success", "diff_hash": "abc123"} + # It hit the runs endpoint for the right run, read-only. + assert client.calls[0][0].endswith("/repos/o/r/actions/runs/123") + + +def test_failure_conclusion_is_echoed_not_judged() -> None: + # The fetcher returns the raw conclusion; it does NOT decide pass/fail. + client = _FakeClient(_FakeResponse(200, {"id": 5, "conclusion": "failure"})) + out = fetch_ci_result(_state(run_id="5"), owner="o", repo="r", client=client) + assert out is not None + assert out["conclusion"] == "failure" # echoed verbatim, no verdict derived + + +def test_diff_hash_echoed_from_state() -> None: + client = _FakeClient(_FakeResponse(200, {"id": 9, "conclusion": "success"})) + out = fetch_ci_result( + {"run_id": "9", "diff_hash": "DEADBEEF"}, owner="o", repo="r", client=client + ) + assert out["diff_hash"] == "DEADBEEF" + + +def test_run_id_read_from_state_drives_url() -> None: + client = _FakeClient(_FakeResponse(200, {"id": 777, "conclusion": "success"})) + fetch_ci_result(_state(run_id="777"), owner="acme", repo="svc", client=client) + assert client.calls[0][0].endswith("/repos/acme/svc/actions/runs/777") + + +# --------------------------------------------------------------------------- # +# Fail-closed paths (every error -> None) +# --------------------------------------------------------------------------- # + + +def test_404_fails_closed() -> None: + client = _FakeClient(_FakeResponse(404, {"message": "Not Found"})) + assert fetch_ci_result(_state(), owner="o", repo="r", client=client) is None + + +def test_auth_fail_401_fails_closed() -> None: + client = _FakeClient(_FakeResponse(401, {"message": "Bad credentials"})) + assert fetch_ci_result(_state(), owner="o", repo="r", client=client) is None + + +def test_forbidden_403_fails_closed() -> None: + client = _FakeClient(_FakeResponse(403, {"message": "forbidden"})) + assert fetch_ci_result(_state(), owner="o", repo="r", client=client) is None + + +def test_malformed_body_fails_closed() -> None: + client = _FakeClient(_FakeResponse(200, ValueError("not json"))) + assert fetch_ci_result(_state(), owner="o", repo="r", client=client) is None + + +def test_non_object_body_fails_closed() -> None: + client = _FakeClient(_FakeResponse(200, ["not", "a", "dict"])) + assert fetch_ci_result(_state(), owner="o", repo="r", client=client) is None + + +def test_missing_conclusion_fails_closed() -> None: + # In-progress run: conclusion is None -> not an authenticated verdict. + client = _FakeClient(_FakeResponse(200, {"id": 1, "conclusion": None})) + assert ( + fetch_ci_result(_state(run_id="1"), owner="o", repo="r", client=client) is None + ) + + +def test_timeout_fails_closed() -> None: + client = _FakeClient(raise_exc=TimeoutError("timed out")) + assert fetch_ci_result(_state(), owner="o", repo="r", client=client) is None + + +def test_connection_error_fails_closed() -> None: + client = _FakeClient(raise_exc=ConnectionError("refused")) + assert fetch_ci_result(_state(), owner="o", repo="r", client=client) is None + + +def test_missing_run_id_fails_closed() -> None: + client = _FakeClient(_FakeResponse(200, {"id": 1, "conclusion": "success"})) + assert ( + fetch_ci_result({"diff_hash": "x"}, owner="o", repo="r", client=client) is None + ) + # And it never even made a request. + assert client.calls == [] + + +def test_empty_run_id_fails_closed() -> None: + client = _FakeClient(_FakeResponse(200, {"id": 1, "conclusion": "success"})) + assert ( + fetch_ci_result( + {"run_id": "", "diff_hash": "x"}, owner="o", repo="r", client=client + ) + is None + ) + assert client.calls == [] + + +# --------------------------------------------------------------------------- # +# Missing token (no injected client) -> fails closed (does NOT raise) +# --------------------------------------------------------------------------- # + + +def test_missing_token_fails_closed(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.delenv(CI_READ_TOKEN_ENV, raising=False) + monkeypatch.delenv("GITHUB_TOKEN", raising=False) + # No client injected -> it must resolve a token; none set -> None (no raise). + assert fetch_ci_result(_state(), owner="o", repo="r") is None + + +# --------------------------------------------------------------------------- # +# Read-only contract: never derives a verdict, never writes +# --------------------------------------------------------------------------- # + + +def test_fetcher_never_returns_a_verdict_field() -> None: + client = _FakeClient(_FakeResponse(200, {"id": 2, "conclusion": "failure"})) + out = fetch_ci_result(_state(run_id="2"), owner="o", repo="r", client=client) + assert out is not None + # The mapping is DATA only — there is no gate_decision / passed / verdict. + assert set(out.keys()) == {"run_id", "conclusion", "diff_hash"} + assert "passed" not in out + assert "gate_decision" not in out + + +def test_build_ci_result_fetcher_returns_state_callable() -> None: + client = _FakeClient(_FakeResponse(200, {"id": 3, "conclusion": "success"})) + fetcher = build_ci_result_fetcher(owner="o", repo="r", client=client) + out = fetcher({"run_id": "3", "diff_hash": "h"}) + assert out == {"run_id": "3", "conclusion": "success", "diff_hash": "h"} + + +def test_build_ci_result_fetcher_fails_closed_on_error() -> None: + client = _FakeClient(_FakeResponse(500, {"message": "boom"})) + fetcher = build_ci_result_fetcher(owner="o", repo="r", client=client) + assert fetcher({"run_id": "3", "diff_hash": "h"}) is None