feat(agent-team): bind planner/review/builder/verifier nodes to their models

review_loop_llm -> GPT-4.1 cross_reviewer (orchestrator run.py); builders_llm -> DeepSeek fast_coder (INERT, proposes diff text only); verifier_llm -> ci_gate is sole PASS authority, Claude is fix-proposer only. Hardens review_loop.parse_verdict to word-boundary matching, adds a fail-closed subprocess timeout, and bind_review_node (single-arg, no LangGraph config injection). All fail safe on untrusted model output.
This commit is contained in:
Adam Moussa 2026-06-18 12:56:42 -04:00
parent 270ce93b2a
commit 4b17e8ebd4
8 changed files with 2277 additions and 21 deletions

View file

@ -0,0 +1,374 @@
"""DeepSeek-backed builders binding — the real §3.3 / §7.1 P3 build seam.
:mod:`agent_team.nodes.builders` owns the Plane-2 builders *node* (the §3.3.2
box-side trust-control-surface denylist + diff-integrity hash) but deliberately
injects the diff-synthesis step behind an ``DiffBuilder`` seam so the leaf stays
pure and unit-testable. Its committed default
(:func:`agent_team.nodes.builders.default_diff_builder`) is a Claude-billing
*stub* whose docstring (builders.py line ~164) notes the REAL implementation
wires "DeepSeek (mechanical edits, via the local orchestrator)". This module is
that real implementation.
Per the locked design, builders are P3 mechanical edits and route to the
orchestrator's ``fast_coder`` (DeepSeek), NOT to Claude. This module therefore
does NOT call :func:`agent_team.billing.claude_invoke`; it calls the local
orchestrator's ``fast_coder`` to produce the candidate diff.
================================ SECURITY BOUNDARY ========================
Builders are P3 in the locked design and HARD-GATED: the live CI apply/verify
trust boundary (§3.3.2) must clear ``/sh-security-review`` + GPT-4.1 cross-review
BEFORE it goes live. This module is MODEL LOGIC ONLY and MUST stay INERT:
* It PROPOSES a candidate diff as DATA (a :class:`CandidateDiff` record). It
NEVER applies a patch, NEVER shells out to ``git``, NEVER writes to or
otherwise mutates the working tree / filesystem, and NEVER makes a live CI
call. Applying a diff is the GATED CI path — not this module's job.
* The only subprocess this module spawns is a read-only call to the local
orchestrator's ``run.py`` to ask ``fast_coder`` for diff TEXT. That
subprocess is a model invocation, not a patch application: its stdout is
parsed as untrusted data and returned; it touches nothing in the target
repo. There is no ``git apply``/``patch``/``git``/``write_text``/``open(...,
"w")`` path anywhere in this file — by construction, the builder cannot
mutate state.
Because the model output is UNTRUSTED, parsing is defensive and FAILS SAFE: on
unparseable output, an empty/whitespace diff, or any build error, the builder
returns an EMPTY/NO-OP candidate marked ``failed`` (``ok is False``) so the
downstream verifier / CI REJECTS it. It NEVER fabricates a "success" diff.
============================================================================
Wiring note (no node edit): this module is a standalone real binding. The
node's injection point is its ``DiffBuilder`` seam — the coordinator should bind
:func:`default_build` (adapted via :func:`as_diff_builder`) into
:func:`agent_team.nodes.builders.builders_node` / ``build_candidate_diff`` at
startup. That wiring edit is deliberately left to the coordinator; this module
does not edit the node.
"""
from __future__ import annotations
import subprocess
from collections.abc import Callable, Mapping
from dataclasses import dataclass
from pathlib import Path
from typing import Any
from agent_team.state_store import compute_content_hash
__all__ = [
"BuildCallable",
"CandidateDiff",
"as_diff_builder",
"build_candidate_diff",
"default_build",
]
# The injectable build seam: given the rendered build instruction (a string),
# return the model's raw candidate-diff text. Tests pass a fake; the default
# (:func:`default_build`) routes to the orchestrator's DeepSeek ``fast_coder``.
BuildCallable = Callable[[str], str]
# Default subprocess timeout (seconds) for the orchestrator fast_coder call.
_DEFAULT_TIMEOUT_S = 600
@dataclass
class CandidateDiff:
"""A proposed candidate diff emitted as DATA (never applied here).
This is the record the builders pipeline carries downstream. It mirrors the
fields :func:`agent_team.nodes.builders.build_candidate_diff` records on the
task — the unified-diff text plus its content-hash — and adds the explicit
fail-safe flags so an unparseable/failed build is propagated as a NO-OP the
verifier/CI rejects, rather than as a fabricated success.
Attributes:
diff: The candidate unified diff (empty string on a failed/no-op build).
diff_hash: Content hash of ``diff`` via
:func:`agent_team.state_store.compute_content_hash` (always computed,
including over the empty diff, so CI keys against it deterministically).
ok: ``True`` only when a non-empty, plausibly-unified diff was produced.
failed: ``True`` when the build failed or produced nothing usable (the
inverse of :attr:`ok`); kept explicit so a downstream check can read
either flag.
reason: Human-readable explanation when :attr:`failed`; empty when ``ok``.
"""
diff: str
diff_hash: str
ok: bool
failed: bool
reason: str = ""
@classmethod
def success(cls, diff: str) -> "CandidateDiff":
"""Build an ``ok`` candidate from a validated non-empty diff string."""
return cls(
diff=diff,
diff_hash=compute_content_hash(diff.encode("utf-8")),
ok=True,
failed=False,
reason="",
)
@classmethod
def no_op(cls, reason: str) -> "CandidateDiff":
"""Build a FAILED no-op candidate (empty diff) the verifier/CI rejects.
The empty diff is still hashed so the record shape is uniform and CI's
hash check has a deterministic value to compare; the ``failed`` flag is
what makes the downstream reject it.
"""
return cls(
diff="",
diff_hash=compute_content_hash(b""),
ok=False,
failed=True,
reason=reason,
)
@dataclass
class _OrchestratorRoute:
"""Resolved location + runner for the local orchestrator ``run.py``.
Kept as a tiny dataclass (rather than module-level constants) so the default
build call resolves the orchestrator root lazily and a test could swap the
runner without importing the orchestrator. No orchestrator code is imported
at module top (mirrors :func:`agent_team.graph.build_sqlite_checkpointer`'s
deferred-import discipline).
"""
root: Path
timeout_s: int = _DEFAULT_TIMEOUT_S
def _orchestrator_root() -> Path:
"""Resolve the orchestrator root (the dir holding ``run.py``).
This file lives at ``<root>/agent-team/agent_team/nodes/builders_llm.py``,
so the orchestrator root is ``parents[3]`` (nodes -> agent_team -> agent-team
-> <root>). Verified against the real tree: ``parents[2]`` is ``agent-team``,
not the root.
"""
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``.
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`.
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.
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.
"""
route = (
route if route is not None else _OrchestratorRoute(root=_orchestrator_root())
)
run_py = route.root / "run.py"
completed = subprocess.run(
["python3", str(run_py), instruction],
capture_output=True,
text=True,
timeout=route.timeout_s,
check=True,
cwd=str(route.root),
)
return completed.stdout
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.
Pure string assembly over the plan/state (no I/O) so the instruction shape is
directly unit-testable. The instruction tells the coder to emit ONLY a single
unified diff and to stay inside the declared scope — the box-side denylist in
:mod:`agent_team.nodes.builders` is the real enforcement, but reinforcing it
in the prompt keeps the model on-task.
"""
title = str(plan.get("title") or plan.get("task") or "(untitled task)")
scope = plan.get("scope") or []
phases = plan.get("phases") or []
repo = ""
raw_repo = state.get("repo") if isinstance(state, Mapping) else None
if isinstance(raw_repo, str) and raw_repo.strip():
repo = raw_repo.strip()
scope_lines = "\n".join(f" - {p}" for p in scope) or " (no scope declared)"
phase_lines = (
"\n".join(f" {i + 1}. {p}" for i, p in enumerate(phases)) or " (none)"
)
sections = [
(
"You are performing a mechanical code edit. Implement the approved "
"plan below as a SINGLE unified diff in git format. Output ONLY the "
"diff — no prose, no explanation, no code fences. Touch ONLY files "
"within the declared scope. Do NOT modify CI workflows, IAM/policy "
"IaC, branch-protection, CODEOWNERS, or Dependabot config."
),
"",
f"Title: {title}",
]
if repo:
sections += [f"Repository: {repo}"]
sections += [
f"Declared scope (paths you may edit):\n{scope_lines}",
f"Phases:\n{phase_lines}",
]
return "\n".join(sections)
# A line is plausibly part of a unified diff if it opens a git/file/hunk header.
# Used only to validate that the model returned a diff (not prose) and to strip
# the orchestrator's framing lines (e.g. ``[retrieved: ...]``, ``[fast_coder]``)
# that run.py prints before the result body. This is validation/extraction over
# UNTRUSTED text — never application.
_DIFF_HEADER_PREFIXES = (
"diff --git ",
"--- ",
"+++ ",
"@@ ",
"index ",
"rename from ",
"rename to ",
"copy from ",
"copy to ",
"new file mode ",
"deleted file mode ",
"old mode ",
"new mode ",
)
def _extract_diff(text: str) -> str | None:
"""Extract a unified diff from UNTRUSTED model/orchestrator output, or ``None``.
The orchestrator's ``run.py`` prints framing lines (``[retrieved: ...]``, a
``[route]`` line, a blank line) before the agent's result. We locate the
first real diff header (``diff --git`` / ``--- `` / ``@@ ``) and return from
there to the end, stripping a trailing code-fence if the model wrapped the
diff. Returns ``None`` when no diff header is present at all (prose-only /
empty output) so the caller fails SAFE to a no-op candidate. Pure text
inspection — it never executes or applies the diff.
"""
if not isinstance(text, str) or not text.strip():
return None
lines = text.splitlines()
start: int | None = None
for idx, line in enumerate(lines):
stripped = line.strip()
# ``diff --git`` and a real ``--- a/...`` header are the strongest
# signals; a lone ``@@`` hunk header also anchors a body-only diff.
if (
stripped.startswith("diff --git ")
or line.startswith("--- ")
or stripped.startswith("@@ ")
):
start = idx
break
if start is None:
return None
body_lines = lines[start:]
# Drop a trailing markdown fence if the model wrapped the diff in ```.
while body_lines and body_lines[-1].strip() in ("```", ""):
if body_lines[-1].strip() == "```":
body_lines.pop()
break
body_lines.pop()
diff = "\n".join(body_lines).strip()
if not diff:
return None
# Require at least one recognizable diff header line, so a stray ``--- ``
# inside prose cannot masquerade as a diff.
if not any(
any(ln.startswith(p) or ln.strip().startswith(p) for p in _DIFF_HEADER_PREFIXES)
for ln in diff.splitlines()
):
return None
return diff
def build_candidate_diff(
plan: Mapping[str, Any],
state: Mapping[str, Any] | None = None,
*,
build: BuildCallable | None = None,
) -> CandidateDiff:
"""Propose a candidate diff for ``plan`` via DeepSeek ``fast_coder`` (P3).
Renders the approved ``plan`` (+ optional ``state``) into a mechanical-edit
instruction, calls the injected ``build`` callable (default
:func:`default_build`, which routes to the orchestrator's DeepSeek
``fast_coder``), defensively parses the UNTRUSTED result, and returns a
:class:`CandidateDiff` record.
FAIL SAFE (never fabricate success): if ``plan`` is not a mapping, the build
raises, or the output does not parse to a non-empty unified diff, this
returns ``CandidateDiff.no_op(reason)`` — an empty diff marked ``failed`` so
the verifier / CI rejects it. A valid diff yields ``CandidateDiff.success``.
INERT: this function only PROPOSES a diff as data. It does not apply it, run
``git``, or write to the filesystem; applying is the gated CI path.
"""
if not isinstance(plan, Mapping):
return CandidateDiff.no_op("approved plan must be a mapping")
instruction = _render_build_instruction(plan, state or {})
build_fn: BuildCallable = build if build is not None else default_build
try:
raw = build_fn(instruction)
except subprocess.TimeoutExpired:
return CandidateDiff.no_op("build timed out")
except subprocess.CalledProcessError as exc:
return CandidateDiff.no_op(f"build process failed (exit {exc.returncode})")
except Exception as exc: # noqa: BLE001 - any builder failure must fail SAFE
return CandidateDiff.no_op(f"build error: {type(exc).__name__}")
if not isinstance(raw, str) or not raw.strip():
return CandidateDiff.no_op("builder produced empty output")
diff = _extract_diff(raw)
if diff is None:
return CandidateDiff.no_op("builder output is not a usable unified diff")
return CandidateDiff.success(diff)
def as_diff_builder(
build: BuildCallable | None = None,
) -> Callable[..., str]:
"""Adapt this binding to the node's ``DiffBuilder`` seam (keyword signature).
:func:`agent_team.nodes.builders.build_candidate_diff` calls its injected
``DiffBuilder`` as ``builder(plan=..., config=...)`` and expects a unified
-diff STRING back (it then hashes + denylist-scans). This adapter lets the
coordinator bind the real DeepSeek path there: it runs
:func:`build_candidate_diff` and returns the diff string on success.
On a failed/no-op build it returns an EMPTY string. The node treats an empty
diff as ``BuildError`` (its own fail-closed contract), so the adapter never
smuggles a fabricated success past the node either. (The richer
:class:`CandidateDiff` record path is available directly via
:func:`build_candidate_diff` for callers that want the explicit failed flag.)
"""
def _builder(*, plan: Mapping[str, Any], config: Mapping[str, Any] | None) -> str:
state = config if isinstance(config, Mapping) else {}
candidate = build_candidate_diff(plan, state, build=build)
return candidate.diff
return _builder

View file

@ -43,6 +43,7 @@ from __future__ import annotations
import json
import os
import re
import subprocess
from dataclasses import dataclass
from datetime import datetime, timezone
@ -57,6 +58,7 @@ __all__ = [
"ReviewOutcome",
"ReviewResult",
"ReviewVerdict",
"bind_review_node",
"review_node",
"route_after_review",
"set_review_invoker",
@ -82,21 +84,59 @@ _MAX_ROUNDS_ENV = "AGENT_TEAM_MAX_REVIEW_ROUNDS"
_RUN_PY_CONFIG_KEY = "orchestrator_run_py"
_DEFAULT_RUN_PY = os.path.expanduser("~/Documents/repositories/orchestrator/run.py")
# Verdict tokens the reviewer output is scanned for. REQUEST_CHANGES wins on a
# tie so an ambiguous review fails closed (loops back / escalates) rather than
# advancing a plan the reviewer flagged.
_APPROVE_TOKENS = ("APPROVE", "APPROVED", "LGTM", "NO BLOCKERS", "NO BLOCKING")
# Config / env keys + default for the orchestrator subprocess timeout (seconds).
# Bounds the default shell-out so a hung run.py cannot stall the bounded review
# loop. Mirrors the resolver in review_loop_llm so the two cannot drift.
_TIMEOUT_CONFIG_KEY = "review_timeout_seconds"
_TIMEOUT_ENV = "AGENT_TEAM_REVIEW_TIMEOUT_SECONDS"
_DEFAULT_TIMEOUT_SECONDS = 600.0
# Sentinel returned by the default invoker when the orchestrator subprocess
# times out. It parses (via parse_verdict) to REQUEST_CHANGES, so a hung run.py
# fails CLOSED (loops back / escalates) instead of blocking the bounded loop.
_TIMEOUT_VERDICT_TEXT = (
"REQUEST CHANGES: orchestrator review timed out (failing closed)."
)
# Verdict tokens the reviewer output is scanned for, matched on WORD BOUNDARIES
# (not substrings). The change marker is the sh-plan-review rubric token
# ``BLOCK`` (e.g. "BLOCK: ..."), matched as a whole word so it does NOT fire
# inside benign prose like "no blockers" / "no blocking issues" (the substring
# false-positive this guards against — those inflected words are deliberately
# NOT change tokens). The former "NO BLOCKERS"/"NO BLOCKING" approve tokens
# existed only to undo that substring false-positive; with word-boundary
# matching they are unreachable (such bare prose is genuinely ambiguous and must
# fail closed), so they are dropped. REQUEST_CHANGES still wins on a tie so an
# ambiguous review fails closed (loops back / escalates) rather than advancing a
# flagged plan.
_APPROVE_TOKENS = ("APPROVE", "APPROVED", "LGTM")
_CHANGES_TOKENS = (
"REQUEST CHANGES",
"REQUEST_CHANGES",
"REQUESTCHANGES",
"BLOCK",
"BLOCKING",
"NEEDS CHANGES",
"NEEDS WORK",
)
def _compile_token_pattern(tokens: tuple[str, ...]) -> re.Pattern[str]:
"""Compile an alternation of ``tokens`` matched on word boundaries.
Word-boundary anchoring is what keeps ``BLOCK`` from matching inside
``BLOCKERS``/``BLOCKING`` (the substring false-positive this guards against).
Tokens are sorted longest-first so a multi-word token (e.g. ``REQUEST
CHANGES``) is preferred over a shorter overlapping one.
"""
ordered = sorted(tokens, key=len, reverse=True)
alternation = "|".join(re.escape(tok) for tok in ordered)
return re.compile(rf"\b(?:{alternation})\b", re.IGNORECASE)
_APPROVE_RE = _compile_token_pattern(_APPROVE_TOKENS)
_CHANGES_RE = _compile_token_pattern(_CHANGES_TOKENS)
class ReviewVerdict(Enum):
"""The adversarial reviewer's verdict on a plan (design §3.3)."""
@ -155,7 +195,13 @@ class ReviewResult:
ReviewInvoker = Callable[..., str]
def _orchestrator_invoker(prompt: str, *, run_py: str, **_kw: Any) -> str:
def _orchestrator_invoker(
prompt: str,
*,
run_py: str,
config: Mapping[str, Any] | None = None,
**_kw: Any,
) -> str:
"""Default reviewer: call the orchestrator's ``cross_reviewer`` (GPT-4.1).
Invokes the local ``run.py`` with the review prompt. The orchestrator's
@ -163,18 +209,28 @@ def _orchestrator_invoker(prompt: str, *, run_py: str, **_kw: Any) -> str:
keeps the review cross-family (a different model than the Claude planner)
and API-billed + LangSmith-traced per design §3.2. Returns the orchestrator's
stdout (the reviewer's verdict + findings).
The call is bounded by :func:`_resolve_timeout`. If ``run.py`` hangs past the
timeout the subprocess is killed and a REQUEST_CHANGES sentinel is returned
(fail CLOSED) so a stuck review cannot block the bounded loop — rather than
raising, which would crash :func:`review_node` (it does not wrap the call).
"""
if not os.path.exists(run_py):
raise FileNotFoundError(
f"orchestrator entry point not found: {run_py}; set "
f"config[{_RUN_PY_CONFIG_KEY!r}] or rebind via set_review_invoker()."
)
completed = subprocess.run( # noqa: S603 - args are not shell-interpolated
["python3", run_py, prompt],
capture_output=True,
text=True,
check=False,
)
try:
completed = subprocess.run( # noqa: S603 - args are not shell-interpolated
["python3", run_py, prompt],
capture_output=True,
text=True,
check=False,
timeout=_resolve_timeout(config),
)
except subprocess.TimeoutExpired:
# Hung run.py -> fail CLOSED (do not block the bounded loop).
return _TIMEOUT_VERDICT_TEXT
if completed.returncode != 0:
raise RuntimeError(
"orchestrator review call failed "
@ -226,6 +282,27 @@ def _resolve_max_rounds(config: Mapping[str, Any] | None) -> int:
return value
def _resolve_timeout(config: Mapping[str, Any] | None) -> float:
"""Resolve the orchestrator subprocess timeout (seconds) from config/env/default.
A non-positive or non-numeric value falls back to the default so a
misconfigured knob cannot disable the bound. Mirrors the resolver in
:mod:`agent_team.nodes.review_loop_llm` so the two cannot drift.
"""
raw: Any = None
if config is not None:
raw = config.get(_TIMEOUT_CONFIG_KEY)
if raw is None:
raw = os.environ.get(_TIMEOUT_ENV)
if raw is None:
return _DEFAULT_TIMEOUT_SECONDS
try:
value = float(raw)
except (TypeError, ValueError):
return _DEFAULT_TIMEOUT_SECONDS
return value if value > 0 else _DEFAULT_TIMEOUT_SECONDS
def _resolve_run_py(config: Mapping[str, Any] | None) -> str:
"""Resolve the orchestrator ``run.py`` path from config, env, or default."""
if config is not None:
@ -242,17 +319,17 @@ def parse_verdict(text: str) -> ReviewVerdict:
"""Parse a :class:`ReviewVerdict` from the reviewer's free text.
Scans for explicit ``REQUEST CHANGES`` / ``BLOCK`` tokens and ``APPROVE`` /
``LGTM`` tokens (case-insensitive). The result **fails closed**: if a
change-requesting token is present, or if neither token class is present
(an ambiguous / empty review), the verdict is ``REQUEST_CHANGES`` so an
unclear review never silently advances a plan to the builders.
``LGTM`` tokens (case-insensitive, **word-boundary** matched so that prose
like "no blocking issues" inside an APPROVE does not trip a change token).
The result **fails closed**: if a change-requesting token is present, or if
neither token class is present (an ambiguous / empty review), the verdict is
``REQUEST_CHANGES`` so an unclear review never silently advances a plan to
the builders.
"""
haystack = (text or "").upper()
has_changes = any(token in haystack for token in _CHANGES_TOKENS)
has_approve = any(token in haystack for token in _APPROVE_TOKENS)
if has_changes:
haystack = text or ""
if _CHANGES_RE.search(haystack):
return ReviewVerdict.REQUEST_CHANGES
if has_approve:
if _APPROVE_RE.search(haystack):
return ReviewVerdict.APPROVE
# Ambiguous / empty review -> fail closed.
return ReviewVerdict.REQUEST_CHANGES
@ -377,6 +454,28 @@ def review_node(
return update
def bind_review_node(
config: Mapping[str, Any] | None = None,
) -> Callable[[PipelineState], PipelineState]:
"""Return a **single-argument** review node bound to ``config`` (P2 wiring).
:func:`review_node` takes an optional ``config`` second argument. If it is
handed to LangGraph directly, LangGraph sees the ``config`` parameter and
injects its own ``RunnableConfig`` there, which (a) emits a typing
``UserWarning`` and (b) means the task's ``max_review_rounds`` / timeout /
``run_py`` overrides never reach the node. Wrapping it as a one-arg closure
over the intended ``config`` keeps the node free of a LangGraph-managed
``config`` param (no warning, no injection) and threads the *real* review
config through. The coordinator passes the bound node to
:func:`agent_team.graph.build_graph` as ``review_node``.
"""
def node(state: PipelineState) -> PipelineState:
return review_node(state, config)
return node
def route_after_review(state: PipelineState) -> str:
"""LangGraph conditional-edge: next node after the review loop.

View file

@ -0,0 +1,292 @@
"""GPT-4.1 cross-family review binding for the REVIEW_LOOP node (design §3.3, §7.1 P2).
:mod:`agent_team.nodes.review_loop` owns the *loop* — the bounded
approve / loop-back-to-planner / escalate-to-Adam state machine — but it
deliberately routes the actual review call through an injectable
:data:`~agent_team.nodes.review_loop.ReviewInvoker` seam (rebound via
``set_review_invoker``) so the loop stays pure and unit-testable.
This module supplies the **real** implementation of that seam: a single
function, :func:`review_plan`, that takes the planner's phased plan plus the
graph state and returns the node-contract verdict
(:class:`~agent_team.nodes.review_loop.ReviewVerdict`) the loop consumes.
Model routing — this node is GPT-4.1, NOT Claude
------------------------------------------------
Per the locked design the review loop runs an **independent, cross-family**
review of the plan via the orchestrator's ``cross_reviewer`` agent
(``gpt-4.1`` — a different model family than the Claude planner, so it catches
different blind spots). The R720 coordinator "reuses the local rsync'd
``orchestrator/run.py`` in place for non-Claude single-shots" (design §3.2), so
the default review call is **not** routed through
:func:`agent_team.billing.claude_invoke` (the *Claude* seam). It instead shells
out to ``python3 <orchestrator_root>/run.py "<review task>"``, whose router
sends adversarial-review tasks to ``cross_reviewer`` (GPT-4.1), keeping the call
API-billed and LangSmith-traced.
``<orchestrator_root>`` resolves to ``Path(__file__).resolve().parents[3]``
(this file lives at ``<root>/agent-team/agent_team/nodes/review_loop_llm.py``,
so parents[0]=nodes, [1]=agent_team, [2]=agent-team, [3]=<root>). The
orchestrator package is **never imported at module load** — the subprocess path
needs no import at all, and the optional in-process fallback
(:func:`make_cross_reviewer_invoker`) imports ``models`` lazily, inside the
call, mirroring the deferred-import discipline of
``graph.build_sqlite_checkpointer``.
The review callable is INJECTABLE (``review=`` parameter / the
``PlanReviewer`` type) so tests pass a fake and never touch the network or a
subprocess. The default is :func:`default_plan_reviewer`.
Contract
--------
* **Input** to the review callable: a single string — the composed review task
(the plan plus prior-round context), built by
:func:`~agent_team.nodes.review_loop.build_review_prompt`.
* **Output** of the review callable: the reviewer's verdict text (free-form),
which :func:`~agent_team.nodes.review_loop.parse_verdict` maps to a
:class:`~agent_team.nodes.review_loop.ReviewVerdict`.
Fail-safe (UNTRUSTED output, never auto-approve on a bad review)
----------------------------------------------------------------
The reviewer's text is untrusted. Parsing is defensive and **fails closed**: if
the call errors (subprocess failure, timeout, exception) or the output is empty
/ unparseable / ambiguous, :func:`review_plan` returns
:attr:`~agent_team.nodes.review_loop.ReviewVerdict.REQUEST_CHANGES` — the SAFE
branch the loop treats as "do not auto-approve" (loop back to the planner, or
escalate to Adam once the round cap is hit). A plan is **only ever** approved on
an explicit, cleanly-parsed APPROVE verdict.
Wiring (one-line injection point, added later — this module edits nothing)
--------------------------------------------------------------------------
At startup a leaf binds this real reviewer into the loop's seam by passing
:func:`default_plan_reviewer` — which already matches the ``ReviewInvoker``
``(prompt, **kw) -> str`` shape (it returns the reviewer's verdict *text*, which
the loop's own :func:`~agent_team.nodes.review_loop.parse_verdict` maps to a
verdict)::
from agent_team.nodes import review_loop
from agent_team.nodes.review_loop_llm import default_plan_reviewer
review_loop.set_review_invoker(default_plan_reviewer)
This is a single line in the wiring module; this file does not edit the node.
"""
from __future__ import annotations
import os
import subprocess
from collections.abc import Callable
from pathlib import Path
from typing import Any
from agent_team.nodes.review_loop import (
ReviewVerdict,
build_review_prompt,
parse_verdict,
)
from agent_team.task_model import PipelineState
__all__ = [
"PlanReviewer",
"default_plan_reviewer",
"make_cross_reviewer_invoker",
"make_run_py_invoker",
"resolve_orchestrator_root",
"resolve_run_py",
"review_plan",
]
# The injectable review callable: takes the composed review task (a string) and
# returns the reviewer's verdict text (a string). The default shells out to the
# orchestrator's cross_reviewer (GPT-4.1) via the local run.py. Tests rebind it
# by passing review=<fake> to review_plan().
PlanReviewer = Callable[..., str]
# Config / env keys naming the orchestrator entry point (the local run.py) and a
# subprocess timeout. Overridable for tests and non-default installs without
# editing this module.
_RUN_PY_CONFIG_KEY = "orchestrator_run_py"
_RUN_PY_ENV = "AGENT_TEAM_ORCHESTRATOR_RUN_PY"
_TIMEOUT_CONFIG_KEY = "review_timeout_seconds"
_TIMEOUT_ENV = "AGENT_TEAM_REVIEW_TIMEOUT_SECONDS"
_DEFAULT_TIMEOUT_SECONDS = 600.0
def resolve_orchestrator_root() -> Path:
"""Return the orchestrator root dir (the parent that holds ``run.py``).
This file lives at ``<root>/agent-team/agent_team/nodes/review_loop_llm.py``,
so the root is ``parents[3]`` of the resolved module path
(parents[0]=nodes, [1]=agent_team, [2]=agent-team, [3]=<root>). Resolved
lazily at call time — the orchestrator package itself is never imported here.
"""
return Path(__file__).resolve().parents[3]
def resolve_run_py(config: Any = None) -> str:
"""Resolve the orchestrator ``run.py`` path from config, then env, then default.
The default is ``<orchestrator_root>/run.py`` (the rsync'd path the R720
coordinator reuses in place, design §3.2). Mirrors the resolution order used
by the loop node so the two cannot drift.
"""
if isinstance(config, dict):
configured = config.get(_RUN_PY_CONFIG_KEY)
if configured:
return os.path.expanduser(str(configured))
env = os.environ.get(_RUN_PY_ENV)
if env:
return os.path.expanduser(env)
return str(resolve_orchestrator_root() / "run.py")
def _resolve_timeout(config: Any = None) -> float:
"""Resolve the subprocess timeout (seconds) from config, env, or default."""
raw: Any = None
if isinstance(config, dict):
raw = config.get(_TIMEOUT_CONFIG_KEY)
if raw is None:
raw = os.environ.get(_TIMEOUT_ENV)
if raw is None:
return _DEFAULT_TIMEOUT_SECONDS
try:
value = float(raw)
except (TypeError, ValueError):
return _DEFAULT_TIMEOUT_SECONDS
return value if value > 0 else _DEFAULT_TIMEOUT_SECONDS
def make_run_py_invoker() -> PlanReviewer:
"""Build the default reviewer: shell out to the orchestrator's ``run.py``.
Returns a callable ``(prompt, *, config=None, **kw) -> str`` that runs
``python3 <run_py> "<prompt>"`` and returns its stdout (the reviewer's
verdict text). The orchestrator's router sends adversarial-review tasks to
``cross_reviewer`` (GPT-4.1), keeping the review cross-family and
API-billed + LangSmith-traced (design §3.2). No orchestrator import is
needed for this path at all.
A missing ``run.py``, a non-zero exit, or a timeout raises — the caller
(:func:`review_plan`) turns any such error into the fail-safe
REQUEST_CHANGES verdict, so a broken review never auto-approves a plan.
"""
def _invoke(prompt: str, *, config: Any = None, **_kw: Any) -> str:
run_py = resolve_run_py(config)
if not os.path.exists(run_py):
raise FileNotFoundError(
f"orchestrator entry point not found: {run_py}; set "
f"config[{_RUN_PY_CONFIG_KEY!r}] or pass review=<callable> to "
"review_plan()."
)
completed = subprocess.run( # noqa: S603 - args are not shell-interpolated
["python3", run_py, prompt],
capture_output=True,
text=True,
check=False,
timeout=_resolve_timeout(config),
)
if completed.returncode != 0:
raise RuntimeError(
"orchestrator review call failed "
f"(exit {completed.returncode}): {completed.stderr.strip()}"
)
return completed.stdout
return _invoke
def make_cross_reviewer_invoker() -> PlanReviewer:
"""Build an in-process reviewer that calls ``cross_reviewer`` (GPT-4.1) directly.
Optional alternative to :func:`make_run_py_invoker` for callers that would
rather invoke the model in-process than spawn ``run.py``. The orchestrator's
``models`` module is imported **lazily, inside the call** (never at module
load), mirroring ``graph.build_sqlite_checkpointer``'s deferred-import
pattern, so importing this module never pulls in the orchestrator stack.
Like the subprocess path, any error propagates so :func:`review_plan` can
fail safe to REQUEST_CHANGES.
"""
def _invoke(prompt: str, *, config: Any = None, **_kw: Any) -> str:
# Deferred import: keep the orchestrator package out of module import.
from models import get_cross_reviewer # noqa: PLC0415
reviewer = get_cross_reviewer()
result = reviewer.invoke(prompt)
text = getattr(result, "content", result)
return text if isinstance(text, str) else str(text)
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()
def review_plan(
plan: Any = None,
state: PipelineState | None = None,
*,
prompt: str | None = None,
review: PlanReviewer | None = None,
config: Any = None,
**kw: Any,
) -> ReviewVerdict:
"""Run one cross-family (GPT-4.1) review of ``plan`` and return the verdict.
This is the REAL implementation of the loop's review seam. It composes the
review task, calls the injected ``review`` callable (default
:func:`default_plan_reviewer`, which shells out to the orchestrator's
``cross_reviewer`` via ``run.py``), and maps the reviewer's free-form text to
a :class:`~agent_team.nodes.review_loop.ReviewVerdict` via the loop's own
:func:`~agent_team.nodes.review_loop.parse_verdict`.
Inputs are flexible so this slots in behind either calling convention:
* ``review_plan(plan, state)`` — compose the prompt from ``plan``/``state``
with :func:`~agent_team.nodes.review_loop.build_review_prompt`; or
* ``review_plan(prompt=...)`` — review an already-composed prompt (this is
the shape the ``ReviewInvoker`` seam hands in, so binding it as the
invoker is a one-liner).
FAIL SAFE: the reviewer output is UNTRUSTED. If the ``review`` call raises
(subprocess failure, timeout, any exception) or returns empty / non-string /
unparseable / ambiguous text, this returns
:attr:`~agent_team.nodes.review_loop.ReviewVerdict.REQUEST_CHANGES` — the
SAFE branch the loop treats as "do not auto-approve" (loop back, or escalate
to Adam at the round cap). A plan is approved **only** on an explicit,
cleanly-parsed APPROVE verdict; a failed or garbage review never approves.
"""
reviewer = review if review is not None else default_plan_reviewer
task = prompt
if task is None:
# Compose the review task from the plan + prior-round context. Accept a
# bare plan dict by adapting it into the minimal state shape the prompt
# builder reads, so callers need not hand-build a full PipelineState.
review_state: PipelineState
if isinstance(state, dict):
review_state = state
elif isinstance(plan, dict):
review_state = {"plan": plan} # type: ignore[assignment]
else:
# No usable plan/state to review -> fail safe, never auto-approve.
return ReviewVerdict.REQUEST_CHANGES
task = build_review_prompt(review_state)
try:
raw = reviewer(task, config=config, **kw)
except Exception:
# Any failure in the review call (subprocess error, timeout, bad import)
# -> fail safe. Never auto-approve a plan on a failed review.
return ReviewVerdict.REQUEST_CHANGES
text = raw if isinstance(raw, str) else "" if raw is None else str(raw)
# parse_verdict itself fails closed on empty / ambiguous text, but we route
# everything through it so the verdict tokens stay single-sourced in the
# loop node and the two cannot drift.
return parse_verdict(text)

View file

@ -0,0 +1,494 @@
"""Claude-backed verifier bindings — the real §3.3 / §3.3.2 P3 logic.
:mod:`agent_team.nodes.verifier` owns the VERIFY-stage LangGraph node and the
phase transitions, but it injects its reasoning seam (the fix-advisor,
:data:`~agent_team.nodes.verifier.FixAdvisor`) so the loop stays pure and
testable. This module supplies the **real** implementation of that seam, plus a
thin pure-code verdict wrapper, mirroring how :mod:`clarifier_llm` backs the
clarifier seams.
The single load-bearing rule from design §3.3.2 boundary #4 is enforced
STRUCTURALLY by the shape of this module, not by convention:
**The LLM verifier cannot declare green.** Pass/fail is owned by a pure-code
gate (:mod:`agent_team.ci_gate`) over the authenticated, patch-independent CI
Checks result (keyed to ``run_id`` + ``diff_hash``). The LLM verifier may
PROPOSE fixes but can NEVER flip the verdict to pass.
So this module keeps two things rigorously SEPARATE:
1. :func:`evaluate_verdict` — a PURE-CODE function that maps an authenticated CI
Checks result -> pass/fail. It is a thin compose over
:func:`agent_team.ci_gate.evaluate_ci_gate`; it reuses that gate verbatim and
does NOT reimplement or weaken it. It FAILS SAFE: a missing, ambiguous, or
unauthenticated CI result is never a pass.
2. :class:`ClaudeFixProposer` — an optional, injectable LLM fix-PROPOSER
(Claude, via :func:`agent_team.billing.claude_invoke`). It is consulted ONLY
on a non-pass verdict to author advisory fix hints for the builders. Its
output is advisory DATA only; it is structurally incapable of changing the
verdict because the verdict is computed first, by the pure-code gate, and is
never read back from the proposer.
Both halves meet in :func:`propose_for_failure`, which computes the verdict with
the gate, and ONLY if that verdict is not a pass consults the proposer for a
hint. The pass branch never touches the LLM at all.
INERT / HARD-GATE NOTE (§3.3.2 P3): the verifier is hard-gated behind
``/sh-security-review`` + a GPT-4.1 cross-review of the CI trust boundary before
it goes live. This module authors the LOGIC ONLY and stays inert: it does NO
live CI dispatch, NO network I/O, and NO filesystem mutation. The authenticated
CI result is passed in as data (the caller fetches it via the read-only PAT),
exactly as :func:`agent_team.ci_gate.evaluate_ci_gate` expects. The Claude call
goes only through the committed billing seam and is injectable, so this module
is fully unit-testable with no SDK or network.
Defensive parsing is a hard requirement: the model output is UNTRUSTED. The
proposer parser fails SAFE — a missing or garbled proposal degrades to an empty
advisory hint and NEVER crashes, and (by construction) never affects the
verdict.
"""
from __future__ import annotations
import json
import re
from collections.abc import Callable, Mapping, Sequence
from typing import Any
from agent_team.billing import ClaudeResult, claude_invoke
from agent_team.ci_gate import GateDecision, GateResult, evaluate_ci_gate
__all__ = [
"FixProposal",
"ClaudeFixProposer",
"evaluate_verdict",
"propose_for_failure",
"build_fix_advisor",
]
# The signature the billing seam exposes: ``claude_invoke(prompt, *, mode=None,
# config=None, **kw) -> ClaudeResult``. Injected so tests pass a fake, mirroring
# the injection pattern used across this codebase (billing.set_invoker, the
# clarifier callables, the verifier node's FixAdvisor seam).
ClaudeInvoke = Callable[..., ClaudeResult]
# System framing handed to Claude when authoring a fix hint. Kept as a module
# constant so callers can override via the constructor without forking the class.
_DEFAULT_SYSTEM = (
"You are the VERIFY stage of an agentic SDLC pipeline. A pure-code gate has "
"ALREADY decided this candidate diff did NOT pass CI; that decision is final "
"and is not yours to make or revisit. Your only job is to read the gate's "
"failure reasons and propose concrete, minimal fixes for the builders to "
"try next. You cannot declare the task green; only the authenticated CI "
"gate can."
)
def evaluate_verdict(
*,
candidate_diff: str,
ledger_hash: str | None,
ci_result: Mapping[str, Any] | None,
expected_run_id: str,
allowed_scope: Sequence[str] | None = None,
) -> GateResult:
"""Compute the pass/fail/block verdict from the authenticated CI result.
This is the §3.3.2 boundary #4 pass authority and the ONLY thing in this
module that can produce a :data:`~agent_team.ci_gate.GateDecision.PASS`. It
is a thin compose over :func:`agent_team.ci_gate.evaluate_ci_gate` — it
reuses that committed pure-code gate verbatim and does not reimplement,
relax, or second-guess any of its rules. The LLM is intentionally NOT a
parameter here: the verdict is derived SOLELY from the authenticated,
patch-independent CI Checks result (keyed to ``run_id`` + ``diff_hash``).
FAILS SAFE. Anything other than an unambiguous authenticated success is a
non-pass:
* a missing ``candidate_diff`` -> :data:`GateDecision.BLOCK` (nothing to
verify; refuse to proceed, never pass);
* a missing/``None`` ``ci_result`` -> ``BLOCK`` (no authenticated result;
the gate never passes without one);
* a run-id mismatch, hash mismatch, denylist hit, or ambiguous/unknown CI
conclusion -> ``BLOCK`` (per the gate);
* a recognised CI failure -> :data:`GateDecision.FAIL`;
* an authenticated ``success`` keyed to the expected run -> ``PASS``.
Returns the gate's :class:`~agent_team.ci_gate.GateResult` unchanged so the
decision stays auditable (its ``reasons`` quote the exact CI conclusion
consumed). Raises :class:`~agent_team.ci_gate.CiGateError` only on
structurally invalid inputs, exactly as the underlying gate does.
"""
if not isinstance(candidate_diff, str):
# No diff to verify is itself a refuse-to-proceed (mirrors the verifier
# node): BLOCK rather than declare anything. Never a pass.
return GateResult(
decision=GateDecision.BLOCK,
reasons=["no candidate diff present to verify"],
run_id=expected_run_id if isinstance(expected_run_id, str) else None,
diff_hash=ledger_hash,
ci_conclusion=None,
)
return evaluate_ci_gate(
candidate_diff=candidate_diff,
ledger_hash=ledger_hash,
ci_result=ci_result,
expected_run_id=expected_run_id,
allowed_scope=allowed_scope,
)
class FixProposal:
"""An advisory fix proposal authored by the LLM (DATA, never a verdict).
Carries only suggestions for the builders: a free-text ``hint`` and an
optional ordered list of ``suggestions``. It deliberately has NO notion of
pass/fail and exposes no way to express one — it is impossible to encode a
"this passed" signal here, which is what structurally guarantees the LLM
cannot declare green (§3.3.2 boundary #4). The verdict is computed entirely
separately by :func:`evaluate_verdict`.
"""
__slots__ = ("hint", "suggestions")
def __init__(self, hint: str = "", suggestions: list[str] | None = None) -> None:
self.hint = hint
self.suggestions = list(suggestions) if suggestions else []
def __bool__(self) -> bool:
return bool(self.hint or self.suggestions)
def __eq__(self, other: object) -> bool:
if not isinstance(other, FixProposal):
return NotImplemented
return self.hint == other.hint and self.suggestions == other.suggestions
def __repr__(self) -> str:
return f"FixProposal(hint={self.hint!r}, suggestions={self.suggestions!r})"
def as_hint(self) -> str:
"""Render this proposal as a single advisory hint string for builders."""
parts: list[str] = []
if self.hint:
parts.append(self.hint)
for idx, suggestion in enumerate(self.suggestions, start=1):
parts.append(f"{idx}. {suggestion}")
return "\n".join(parts)
# An empty proposal — the fail-safe result whenever the model is unwired,
# errors, or returns garbage. It changes nothing and asserts nothing.
_EMPTY_PROPOSAL = FixProposal()
class ClaudeFixProposer:
"""Claude-backed fix PROPOSER — advisory only, never a verdict (§3.3.2 P3).
Construct with an optional ``invoke`` callable (defaults to
:func:`agent_team.billing.claude_invoke`) so tests inject a fake and the
real wiring goes through the billing seam. ``model`` / ``config`` are passed
through to the invoker, and ``system`` overrides the prompt framing.
The proposer is consulted ONLY on a non-pass :class:`GateResult` to author a
next-fix hint from the *failure* reasons. It returns a :class:`FixProposal`,
which is pure advisory DATA — it carries no pass/fail and cannot influence
the verdict, which is computed independently by :func:`evaluate_verdict`.
Every failure mode degrades to an empty proposal rather than raising: an
unbound/throwing invoker, a non-string reply, or unparseable JSON all yield
:data:`_EMPTY_PROPOSAL`. A garbage proposal therefore never crashes the
pipeline and never changes the verdict.
"""
def __init__(
self,
*,
invoke: ClaudeInvoke | None = None,
model: str | None = None,
config: Any = None,
system: str = _DEFAULT_SYSTEM,
) -> None:
self._invoke: ClaudeInvoke = invoke if invoke is not None else claude_invoke
self._model = model
self._config = config
self._system = system
def propose(
self,
gate_result: GateResult,
state: Mapping[str, Any] | None = None,
) -> FixProposal:
"""Return an advisory :class:`FixProposal` for a non-pass gate result.
On a :data:`GateDecision.PASS` this returns an empty proposal WITHOUT
calling the model: the LLM is never consulted on success, structurally
keeping it off the happy path. On any other decision it asks Claude for
fix suggestions and parses the reply defensively, failing SAFE to an
empty proposal on any trouble (unbound invoker, non-string reply, bad
JSON). It never raises and never returns anything that could read as a
verdict.
"""
if gate_result.decision is GateDecision.PASS:
# The gate already passed; the LLM has no role here. Never consulted.
return _EMPTY_PROPOSAL
prompt = self._build_prompt(gate_result, state or {})
try:
result = self._invoke(prompt, model=self._model, config=self._config)
text = getattr(result, "text", "")
except Exception:
# An unwired or throwing invoker must not crash VERIFY; the verdict
# already stands and the builders simply loop back without a hint.
return _EMPTY_PROPOSAL
return self._parse(text)
# The proposer matches the verifier node's FixAdvisor seam:
# ``(GateResult, Mapping) -> str``. Returning the rendered hint string keeps
# the advisory output as plain DATA the node appends to its verdict record.
def advise(self, gate_result: GateResult, state: Mapping[str, Any]) -> str:
"""Adapt :meth:`propose` to the verifier node's ``FixAdvisor`` seam.
Matches :data:`agent_team.nodes.verifier.FixAdvisor` exactly
(``(GateResult, Mapping) -> str``) so it wires straight into
:func:`agent_team.nodes.verifier.set_fix_advisor`. Returns the rendered
advisory hint (empty string when there is nothing to suggest) — never a
verdict.
"""
return self.propose(gate_result, state).as_hint()
def _build_prompt(self, gate_result: GateResult, state: Mapping[str, Any]) -> str:
"""Assemble the fix-proposer prompt from the gate failure + task state.
Pure string assembly over the gate result and graph state — no I/O — so
the prompt shape is directly unit-testable.
"""
description = _task_description(state)
reasons = "\n".join(f"- {r}" for r in gate_result.reasons) or "(none recorded)"
ci_conclusion = gate_result.ci_conclusion or "(no authenticated conclusion)"
sections: list[str] = [
self._system,
"",
"## Task",
description or "(no task description provided)",
"",
"## Pure-code gate decision (FINAL, not yours to change)",
f"decision: {gate_result.decision.value}",
f"run_id: {gate_result.run_id}",
f"ci_conclusion: {ci_conclusion}",
"",
"## Gate failure reasons",
reasons,
"",
"## Your job",
(
"Propose the next concrete, minimal fixes for the builders. Do "
"NOT claim the task passed or is green; that verdict is owned by "
"the authenticated CI gate above, not by you."
),
"",
"## Output format",
(
"Respond with ONLY a strict JSON object and no prose outside it, "
'with keys: "hint" (a short string summary) and "suggestions" (a '
"list of strings, most promising first). Example: "
'{"hint": "...", "suggestions": ["...", "..."]}'
),
]
return "\n".join(sections)
def _parse(self, text: str) -> FixProposal:
"""Parse the UNTRUSTED model reply into a :class:`FixProposal`.
Fails SAFE at every step: a non-string reply, no parseable JSON object,
or missing keys all collapse to an empty proposal. Because the proposal
type cannot express a verdict, even a maximally adversarial reply
("everything passed!") cannot influence pass/fail. Never raises.
"""
data = _extract_json_object(text)
if data is None:
return _EMPTY_PROPOSAL
hint = ""
raw_hint = data.get("hint")
if isinstance(raw_hint, str):
hint = raw_hint.strip()
suggestions = _coerce_suggestions(data.get("suggestions"))
if not hint and not suggestions:
return _EMPTY_PROPOSAL
return FixProposal(hint=hint, suggestions=suggestions)
def propose_for_failure(
*,
candidate_diff: str,
ledger_hash: str | None,
ci_result: Mapping[str, Any] | None,
expected_run_id: str,
allowed_scope: Sequence[str] | None = None,
proposer: ClaudeFixProposer | None = None,
state: Mapping[str, Any] | None = None,
) -> tuple[GateResult, FixProposal]:
"""Compute the verdict, then (only on a non-pass) get an advisory proposal.
This is where the two halves meet WITHOUT letting the LLM near the verdict:
1. The verdict is computed FIRST by :func:`evaluate_verdict` (the pure-code
gate). This is the sole pass authority.
2. ONLY if that verdict is not a :data:`GateDecision.PASS` is the
``proposer`` consulted for an advisory :class:`FixProposal`. On a pass,
the proposer is never called and an empty proposal is returned.
The returned ``GateResult`` is exactly what the gate produced — the proposal
is never read back into it — so a garbage or "this passed!" LLM reply cannot
flip a failing verdict to pass. Returns ``(gate_result, proposal)``.
"""
gate_result = evaluate_verdict(
candidate_diff=candidate_diff,
ledger_hash=ledger_hash,
ci_result=ci_result,
expected_run_id=expected_run_id,
allowed_scope=allowed_scope,
)
if gate_result.decision is GateDecision.PASS:
return gate_result, _EMPTY_PROPOSAL
active_proposer = proposer if proposer is not None else ClaudeFixProposer()
proposal = active_proposer.propose(gate_result, state)
return gate_result, proposal
def build_fix_advisor(
*,
invoke: ClaudeInvoke | None = None,
model: str | None = None,
config: Any = None,
system: str = _DEFAULT_SYSTEM,
) -> Callable[[GateResult, Mapping[str, Any]], str]:
"""Build the ``FixAdvisor`` callable for wiring into the verifier node.
Returns the bound :meth:`ClaudeFixProposer.advise` of a shared proposer,
ready to hand to :func:`agent_team.nodes.verifier.set_fix_advisor`. This is
the documented injection point: the verifier node never imports this module
directly — a leaf calls ``set_fix_advisor(build_fix_advisor(...))`` once at
startup, keeping the node dependency-free and structurally guaranteeing the
advisor is only ever consulted on a gate failure (the node never calls it on
a PASS).
"""
proposer = ClaudeFixProposer(
invoke=invoke, model=model, config=config, system=system
)
return proposer.advise
# --------------------------------------------------------------------------- #
# Module-level helpers (pure; no I/O). Mirror clarifier_llm's defensive parsers.
# --------------------------------------------------------------------------- #
def _task_description(state: Mapping[str, Any]) -> str:
"""Pull the task description out of the graph state, defensively.
Looks in the conventional places (the ``plan`` dict, then a top-level
``task``/``description`` key) and falls back to an empty string so a
malformed state surfaces as an empty prompt section, never a ``KeyError``.
"""
plan = state.get("plan") or {}
if isinstance(plan, dict):
desc = plan.get("task") or plan.get("description")
if isinstance(desc, str) and desc.strip():
return desc.strip()
for key in ("task", "description"):
value = state.get(key)
if isinstance(value, str) and value.strip():
return value.strip()
return ""
def _coerce_suggestions(value: Any) -> list[str]:
"""Coerce the model's suggestion list into clean non-empty strings.
Anything that is not a list of usable strings collapses to an empty list.
"""
if not isinstance(value, list):
return []
out: list[str] = []
for item in value:
if isinstance(item, str):
text = item.strip()
if text:
out.append(text)
return out
# A fenced ```json ... ``` block, if the model wrapped its JSON in Markdown.
_FENCE_RE = re.compile(
r"```(?:json)?\s*\n?(?P<body>.*?)\n?\s*```",
flags=re.DOTALL | re.IGNORECASE,
)
def _extract_json_object(text: str) -> dict[str, Any] | None:
"""Extract a JSON object from UNTRUSTED model output, or ``None``.
Tolerates a leading apology or trailing prose and ```json fences. Tries, in
order, the whole string, the contents of a fenced block, then the first
``{...}`` span found by brace matching. Returns ``None`` (never raises) when
nothing parses to a JSON object, so the caller can fail SAFE.
"""
if not isinstance(text, str) or not text.strip():
return None
candidates: list[str] = [text.strip()]
fence = _FENCE_RE.search(text)
if fence:
candidates.append(fence.group("body").strip())
span = _first_brace_span(text)
if span is not None:
candidates.append(span)
for candidate in candidates:
if not candidate:
continue
try:
parsed = json.loads(candidate)
except (json.JSONDecodeError, ValueError):
continue
if isinstance(parsed, dict):
return parsed
return None
def _first_brace_span(text: str) -> str | None:
"""Return the first balanced ``{...}`` span in ``text`` (string-aware)."""
start = text.find("{")
if start == -1:
return None
depth = 0
in_string = False
escaped = False
for idx in range(start, len(text)):
ch = text[idx]
if in_string:
if escaped:
escaped = False
elif ch == "\\":
escaped = True
elif ch == '"':
in_string = False
continue
if ch == '"':
in_string = True
elif ch == "{":
depth += 1
elif ch == "}":
depth -= 1
if depth == 0:
return text[start : idx + 1]
return None

View file

@ -0,0 +1,354 @@
"""Unit tests for agent_team.nodes.builders_llm (§3.3, §7.1 P3).
The DeepSeek-backed builders binding is exercised with a FAKE ``build`` callable
that returns canned text — no network, no subprocess. The load-bearing
properties under test:
* **Clean import.** The module imports without importing the orchestrator at
module top (the orchestrator package is not importable from this tree).
* **Happy path.** A fake build returning a valid unified diff yields an ``ok``
:class:`CandidateDiff` with the right diff and a real content hash.
* **Fail SAFE.** Garbage / empty / prose-only model output yields a FAILED
no-op candidate (empty diff, ``ok is False``), never a fabricated success.
A build that raises also fails safe.
* **Inert boundary.** The module exposes no patch-applying / git / fs-write
function — by construction it cannot mutate the repo.
"""
from __future__ import annotations
import inspect
import sys
from pathlib import Path
from typing import Any
import pytest
from agent_team.nodes import builders, builders_llm
from agent_team.nodes.builders import BuildError, builders_node
from agent_team.nodes.builders_llm import (
CandidateDiff,
as_diff_builder,
build_candidate_diff,
default_build,
)
from agent_team.state_store import compute_content_hash
# A minimal but realistic unified diff the fake build can return.
_VALID_DIFF = (
"diff --git a/agent_team/example.py b/agent_team/example.py\n"
"--- a/agent_team/example.py\n"
"+++ b/agent_team/example.py\n"
"@@ -1,2 +1,2 @@\n"
"-old = 1\n"
"+new = 2\n"
)
_PLAN = {
"title": "Add a thing",
"scope": ["agent_team/"],
"phases": ["edit example.py"],
}
# --------------------------------------------------------------------------- #
# Import hygiene
# --------------------------------------------------------------------------- #
def test_module_imports_without_orchestrator_at_top() -> None:
"""The module imports cleanly with NO orchestrator import at module top.
Parses the module's own top-level import statements (AST) and asserts none of
them pull in the orchestrator's top-level modules — the real DeepSeek route
defers to a subprocess, mirroring graph.build_sqlite_checkpointer's deferred
import. (We inspect the AST rather than reload the module, so the
``CandidateDiff`` identity used by other tests stays stable.)
"""
import ast
tree = ast.parse(inspect.getsource(builders_llm))
orchestrator_mods = {"graph", "agents", "run", "models", "tools", "retriever"}
top_level_imports: set[str] = set()
for node in tree.body: # module body only -> top-level imports
if isinstance(node, ast.Import):
top_level_imports.update(alias.name.split(".")[0] for alias in node.names)
elif isinstance(node, ast.ImportFrom) and node.module:
top_level_imports.add(node.module.split(".")[0])
leaked = top_level_imports & orchestrator_mods
assert not leaked, f"builders_llm imports orchestrator modules at top: {leaked}"
# And the module imports cleanly (already imported above).
assert builders_llm.build_candidate_diff is build_candidate_diff
assert "agent_team.nodes.builders_llm" in sys.modules
# --------------------------------------------------------------------------- #
# Happy path
# --------------------------------------------------------------------------- #
def test_valid_diff_yields_ok_candidate() -> None:
calls: list[str] = []
def fake_build(instruction: str) -> str:
calls.append(instruction)
return _VALID_DIFF
candidate = build_candidate_diff(_PLAN, {"repo": "demo"}, build=fake_build)
assert isinstance(candidate, CandidateDiff)
assert candidate.ok is True
assert candidate.failed is False
assert candidate.reason == ""
assert candidate.diff == _VALID_DIFF.strip()
assert candidate.diff_hash == compute_content_hash(candidate.diff.encode("utf-8"))
# The instruction was rendered from the plan and handed to the builder.
assert calls and "Add a thing" in calls[0]
assert "unified diff" in calls[0]
def test_orchestrator_framing_lines_are_stripped() -> None:
"""run.py prints framing lines before the result; they must be stripped."""
framed = "[retrieved: none]\n[fast_coder]\n\n" + _VALID_DIFF
candidate = build_candidate_diff(_PLAN, build=lambda _i: framed)
assert candidate.ok is True
assert candidate.diff.startswith("diff --git ")
assert "[fast_coder]" not in candidate.diff
assert "[retrieved" not in candidate.diff
def test_fenced_diff_is_unwrapped() -> None:
fenced = "Here is the diff:\n```diff\n" + _VALID_DIFF + "```\n"
candidate = build_candidate_diff(_PLAN, build=lambda _i: fenced)
assert candidate.ok is True
assert candidate.diff.endswith("+new = 2")
assert "```" not in candidate.diff
# --------------------------------------------------------------------------- #
# Fail SAFE — never a fabricated success
# --------------------------------------------------------------------------- #
def test_garbage_output_yields_failed_no_op() -> None:
candidate = build_candidate_diff(
_PLAN, build=lambda _i: "Sure! I cannot produce a diff right now."
)
assert candidate.ok is False
assert candidate.failed is True
assert candidate.diff == ""
assert candidate.reason
assert candidate.diff_hash == compute_content_hash(b"")
def test_empty_output_yields_failed_no_op() -> None:
candidate = build_candidate_diff(_PLAN, build=lambda _i: " \n ")
assert candidate.ok is False
assert candidate.failed is True
assert candidate.diff == ""
def test_build_exception_fails_safe() -> None:
def boom(_instruction: str) -> str:
raise RuntimeError("model exploded")
candidate = build_candidate_diff(_PLAN, build=boom)
assert candidate.ok is False
assert candidate.failed is True
assert candidate.diff == ""
assert "build error" in candidate.reason
def test_non_mapping_plan_fails_safe() -> None:
candidate = build_candidate_diff("not a plan", build=lambda _i: _VALID_DIFF) # type: ignore[arg-type]
assert candidate.ok is False
assert candidate.failed is True
assert candidate.diff == ""
# --------------------------------------------------------------------------- #
# DiffBuilder adapter (node seam parity)
# --------------------------------------------------------------------------- #
def test_as_diff_builder_returns_string_on_success() -> None:
builder = as_diff_builder(build=lambda _i: _VALID_DIFF)
diff = builder(plan=_PLAN, config={"repo": "demo"})
assert isinstance(diff, str)
assert diff.startswith("diff --git ")
def test_as_diff_builder_returns_empty_on_failure() -> None:
"""A failed build must surface as an empty string (node's fail-closed input)."""
builder = as_diff_builder(build=lambda _i: "no diff here")
diff = builder(plan=_PLAN, config=None)
assert diff == ""
# --------------------------------------------------------------------------- #
# Inert / no-apply boundary
# --------------------------------------------------------------------------- #
def test_module_exposes_no_apply_or_fs_mutation_function() -> None:
"""No public callable hints at applying a patch, git, or writing files."""
forbidden_tokens = (
"apply",
"git",
"commit",
"push",
"write",
"mutat",
"patch",
"checkout",
"remove",
"delete",
)
public = [
name
for name in dir(builders_llm)
if not name.startswith("_") and callable(getattr(builders_llm, name))
]
for name in public:
lowered = name.lower()
for token in forbidden_tokens:
assert token not in lowered, (
f"public callable {name!r} suggests a mutation/apply path"
)
def test_source_has_no_patch_application_or_fs_write_paths() -> None:
"""Static guard: NO executable call applies a patch, runs git, or writes files.
Inspects the AST (so the SECURITY-BOUNDARY docstring's mentions of what the
module does NOT do are ignored) and asserts no call/attribute names a
git/patch/apply/fs-mutation primitive. The only subprocess permitted is the
read-only ``subprocess.run`` model call to the orchestrator.
"""
import ast
src = inspect.getsource(builders_llm)
tree = ast.parse(src)
banned_attrs = {
"Popen",
"write_text",
"write_bytes",
"unlink",
"rmtree",
"remove",
"mkdir",
"rename",
"replace",
}
banned_names = {"open"}
subprocess_attrs: set[str] = set()
for node in ast.walk(tree):
if isinstance(node, ast.Attribute):
assert node.attr not in banned_attrs, (
f"module calls a banned fs/git primitive: .{node.attr}"
)
if isinstance(node.value, ast.Name) and node.value.id == "subprocess":
subprocess_attrs.add(node.attr)
if isinstance(node, ast.Name):
assert node.id not in banned_names, (
f"module references a banned builtin: {node.id}"
)
# The only subprocess primitives used are the read-only ``run`` call plus the
# exception types caught around it — never Popen/call/etc. that could shell a
# patch-apply.
assert subprocess_attrs <= {
"run",
"TimeoutExpired",
"CalledProcessError",
}, f"module uses unexpected subprocess primitives: {subprocess_attrs}"
def test_default_build_invokes_run_py_as_list_argv(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""``default_build`` shells ``run.py`` via list-form argv (no shell) + maps stdout.
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.
"""
captured: dict[str, Any] = {}
class _FakeCompleted:
stdout = "DIFF-FROM-SUBPROCESS"
def _fake_run(argv: Any, **kwargs: Any) -> _FakeCompleted:
captured["argv"] = argv
captured["kwargs"] = kwargs
return _FakeCompleted()
monkeypatch.setattr(builders_llm.subprocess, "run", _fake_run)
root = Path("/tmp/fake-orchestrator-root")
route = builders_llm._OrchestratorRoute(root=root)
out = default_build("do the edit", route=route)
assert out == "DIFF-FROM-SUBPROCESS"
# List-form argv (no shell): exactly python3, the run.py path, the instruction.
assert captured["argv"] == ["python3", str(root / "run.py"), "do the edit"]
# No shell=True anywhere in the call (defense against shell injection).
assert captured["kwargs"].get("shell", False) is False
def test_as_diff_builder_empty_raises_build_error_in_real_node() -> None:
"""A failed build surfaces as the real ``builders_node``'s BuildError.
Integration across the seam boundary: ``as_diff_builder`` adapts a failing
build ("no diff here" has no diff header -> empty candidate -> empty string),
and the REAL ``agent_team.nodes.builders.builders_node`` raises its own
``BuildError`` on the empty-diff path rather than emitting a fabricated diff.
"""
builder = as_diff_builder(build=lambda _i: "no diff here")
state = {"plan": dict(_PLAN)}
with pytest.raises(BuildError, match="empty candidate diff"):
builders_node(state, builder=builder) # type: ignore[arg-type]
# Sanity: the node + error are the foundation's, not a redefinition.
assert builders_node.__module__ == builders.__name__
def test_default_build_is_the_injection_default() -> None:
"""The default build seam is default_build (the DeepSeek/orchestrator route)."""
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).
src = inspect.getsource(default_build)
assert "run.py" in src
assert "subprocess.run" in src
# builders are DeepSeek (orchestrator fast_coder), NOT Claude: assert the
# module never CALLS billing.claude_invoke (AST, so docstring mentions of the
# "NOT claude_invoke" contrast don't trip the check).
import ast
tree = ast.parse(inspect.getsource(builders_llm))
called = set()
for node in ast.walk(tree):
if isinstance(node, ast.Call):
fn = node.func
if isinstance(fn, ast.Name):
called.add(fn.id)
elif isinstance(fn, ast.Attribute):
called.add(fn.attr)
assert "claude_invoke" not in called

View file

@ -94,6 +94,35 @@ def test_parse_verdict_request_changes_wins_on_conflict() -> None:
assert parse_verdict(text) is ReviewVerdict.REQUEST_CHANGES
def test_parse_verdict_approve_with_no_blockers_prose() -> None:
# Regression: "no blockers" / "no blocking" prose inside an APPROVE must not
# trip the BLOCK change-token (substring false-positive). Word-boundary
# matching keeps these as APPROVE.
assert parse_verdict("Approved, no blockers.") is ReviewVerdict.APPROVE
assert (
parse_verdict("VERDICT: APPROVE — no blocking issues found")
is ReviewVerdict.APPROVE
)
assert parse_verdict("LGTM, found no blockers") is ReviewVerdict.APPROVE
def test_parse_verdict_real_block_token_requests_changes() -> None:
# A real, standalone BLOCK verdict token (rubric vocabulary) -> REQUEST_CHANGES.
assert parse_verdict("BLOCK: unsafe IAM policy") is ReviewVerdict.REQUEST_CHANGES
assert (
parse_verdict("VERDICT: REQUEST CHANGES\nthis is a BLOCK")
is ReviewVerdict.REQUEST_CHANGES
)
def test_parse_verdict_bare_no_blockers_prose_fails_closed() -> None:
# "no blockers" with NO explicit APPROVE/LGTM token is genuinely ambiguous
# and must fail closed (the dropped NO BLOCKERS approve token is unreachable).
# Note these inflected words ("blockers"/"blocking") are NOT change tokens.
assert parse_verdict("no blockers") is ReviewVerdict.REQUEST_CHANGES
assert parse_verdict("no blocking issues") is ReviewVerdict.REQUEST_CHANGES
# --------------------------------------------------------------------------- #
# build_review_prompt
# --------------------------------------------------------------------------- #
@ -350,3 +379,64 @@ def test_review_result_to_dict_round_trips_fields() -> None:
def test_default_invoker_missing_run_py_raises() -> None:
with pytest.raises(FileNotFoundError):
review_loop._orchestrator_invoker("prompt", run_py="/nonexistent/path/run.py")
def test_default_invoker_passes_timeout_to_subprocess(
monkeypatch: pytest.MonkeyPatch,
) -> None:
# The default shell-out must pass a bounded timeout to subprocess.run.
captured: dict[str, Any] = {}
class _Completed:
returncode = 0
stdout = "VERDICT: APPROVE"
stderr = ""
def _fake_run(args: list[str], **kw: Any) -> _Completed:
captured["kw"] = kw
return _Completed()
monkeypatch.setattr(review_loop.subprocess, "run", _fake_run)
monkeypatch.setattr(review_loop.os.path, "exists", lambda _p: True)
out = review_loop._orchestrator_invoker(
"prompt", run_py="/tmp/run.py", config={"review_timeout_seconds": 12}
)
assert "APPROVE" in out
assert captured["kw"]["timeout"] == 12.0
def test_default_invoker_timeout_fails_closed(
monkeypatch: pytest.MonkeyPatch,
) -> None:
# A hung run.py (TimeoutExpired) must fail CLOSED: return text that parses to
# REQUEST_CHANGES rather than raising and crashing review_node.
import subprocess as _sp
def _raise_timeout(args: list[str], **kw: Any):
raise _sp.TimeoutExpired(cmd=args, timeout=kw.get("timeout", 1))
monkeypatch.setattr(review_loop.subprocess, "run", _raise_timeout)
monkeypatch.setattr(review_loop.os.path, "exists", lambda _p: True)
out = review_loop._orchestrator_invoker("prompt", run_py="/tmp/run.py")
assert parse_verdict(out) is ReviewVerdict.REQUEST_CHANGES
def test_review_node_survives_timeout_fail_closed(
monkeypatch: pytest.MonkeyPatch,
) -> None:
# End-to-end: a hung default invoker makes review_node loop back / escalate,
# never approve, and never raise.
import subprocess as _sp
def _raise_timeout(args: list[str], **kw: Any):
raise _sp.TimeoutExpired(cmd=args, timeout=kw.get("timeout", 1))
monkeypatch.setattr(review_loop.subprocess, "run", _raise_timeout)
monkeypatch.setattr(review_loop.os.path, "exists", lambda _p: True)
# Use the real default invoker (not a test fake).
set_review_invoker(review_loop._orchestrator_invoker)
update = review_node(_state(), config={"max_review_rounds": 3})
assert (
update["review_verdicts"][-1]["verdict"] == ReviewVerdict.REQUEST_CHANGES.value
)
assert update["current_phase"] == Phase.PLAN.value

View file

@ -0,0 +1,229 @@
"""Unit tests for agent_team.nodes.review_loop_llm (§3.3, §7.1 P2).
The real GPT-4.1 cross-family review binding is exercised with a FAKE review
callable that returns canned verdict text — no network, no subprocess. The
load-bearing properties under test:
* **No orchestrator import at module load.** Importing this module must not pull
in the orchestrator package (``models`` / ``graph``); the default reviewer
shells out / imports lazily.
* **Clear approve -> APPROVE.** An injected fake returning an explicit APPROVE
verdict maps to the node-contract ``ReviewVerdict.APPROVE`` (proceed).
* **Changes requested -> REQUEST_CHANGES.** The loop-back / escalate verdict.
* **Fail SAFE.** A review call that raises, or returns garbage / empty / a
non-string, maps to ``REQUEST_CHANGES`` — never an auto-approve.
* **Routing facts.** The default reviewer resolves the orchestrator root at
``parents[3]`` and a ``run.py`` next to it, and is bound as the default.
"""
from __future__ import annotations
import subprocess
import sys
from typing import Any
import pytest
from agent_team.nodes.review_loop import ReviewVerdict
from agent_team.nodes.review_loop_llm import (
default_plan_reviewer,
make_run_py_invoker,
resolve_orchestrator_root,
resolve_run_py,
review_plan,
)
# --------------------------------------------------------------------------- #
# Fakes / helpers
# --------------------------------------------------------------------------- #
class _FakeReview:
"""A fake plan reviewer that returns canned text and records its calls.
``reply`` is the verdict text returned every call. ``raises`` (if set) is
raised instead, to simulate a failed review call.
"""
def __init__(self, reply: Any = "", *, raises: BaseException | None = None) -> None:
self._reply = reply
self._raises = raises
self.calls: list[dict[str, Any]] = []
def __call__(self, prompt: str, **kw: Any) -> Any:
self.calls.append({"prompt": prompt, "kw": kw})
if self._raises is not None:
raise self._raises
return self._reply
_PLAN = {"task": "ship a thing", "phases": [{"name": "P1"}, {"name": "P2"}]}
def _state() -> dict[str, Any]:
return {"plan": _PLAN, "review_verdicts": []}
# --------------------------------------------------------------------------- #
# Module import hygiene
# --------------------------------------------------------------------------- #
def test_module_imports_without_orchestrator() -> None:
"""Importing the module must not import the orchestrator package."""
# The module is already imported at top, but assert the orchestrator stack
# did not get pulled in as a side effect of importing it.
assert "models" not in sys.modules
assert "graph" not in sys.modules
# --------------------------------------------------------------------------- #
# Verdict mapping
# --------------------------------------------------------------------------- #
def test_clear_approve_maps_to_approve() -> None:
"""An injected fake returning a clear APPROVE verdict -> ReviewVerdict.APPROVE."""
fake = _FakeReview("VERDICT: APPROVE\nThe plan is sound and ready to build.")
verdict = review_plan(_PLAN, _state(), review=fake)
assert verdict is ReviewVerdict.APPROVE
def test_changes_requested_maps_to_request_changes() -> None:
"""An injected fake returning changes-requested -> ReviewVerdict.REQUEST_CHANGES."""
fake = _FakeReview("VERDICT: REQUEST CHANGES\nPhase ordering is wrong.")
verdict = review_plan(_PLAN, _state(), review=fake)
assert verdict is ReviewVerdict.REQUEST_CHANGES
def test_block_token_maps_to_request_changes() -> None:
"""A BLOCK verdict (sh-plan-review vocabulary) -> REQUEST_CHANGES."""
fake = _FakeReview("BLOCK: missing rollback phase.")
assert review_plan(_PLAN, _state(), review=fake) is ReviewVerdict.REQUEST_CHANGES
def test_prompt_only_calling_convention() -> None:
"""review_plan(prompt=...) reviews an already-composed prompt (the seam shape)."""
fake = _FakeReview("APPROVE")
verdict = review_plan(prompt="pre-composed review task", review=fake)
assert verdict is ReviewVerdict.APPROVE
assert fake.calls[0]["prompt"] == "pre-composed review task"
def test_prompt_is_composed_from_plan_when_not_supplied() -> None:
"""With no prompt, the plan text is embedded in the composed review task."""
fake = _FakeReview("APPROVE")
review_plan(_PLAN, _state(), review=fake)
sent = fake.calls[0]["prompt"]
assert "ship a thing" in sent
# --------------------------------------------------------------------------- #
# Fail-safe (UNTRUSTED output, never auto-approve)
# --------------------------------------------------------------------------- #
def test_review_call_raising_fails_safe() -> None:
"""A review call that raises -> REQUEST_CHANGES, never an auto-approve."""
fake = _FakeReview(raises=RuntimeError("orchestrator exploded"))
assert review_plan(_PLAN, _state(), review=fake) is ReviewVerdict.REQUEST_CHANGES
def test_garbage_output_fails_safe() -> None:
"""Unparseable / ambiguous reviewer text -> REQUEST_CHANGES."""
fake = _FakeReview("lorem ipsum dolor sit amet, nothing verdict-like here")
assert review_plan(_PLAN, _state(), review=fake) is ReviewVerdict.REQUEST_CHANGES
def test_empty_output_fails_safe() -> None:
"""Empty reviewer output -> REQUEST_CHANGES (fail closed)."""
fake = _FakeReview("")
assert review_plan(_PLAN, _state(), review=fake) is ReviewVerdict.REQUEST_CHANGES
def test_non_string_output_fails_safe() -> None:
"""A non-string (e.g. None) reviewer output never auto-approves."""
fake = _FakeReview(None)
assert review_plan(_PLAN, _state(), review=fake) is ReviewVerdict.REQUEST_CHANGES
def test_conflicting_tokens_fail_closed() -> None:
"""When both APPROVE and REQUEST CHANGES appear, fail closed (changes wins)."""
fake = _FakeReview("APPROVE in spirit but REQUEST CHANGES on phase 2.")
assert review_plan(_PLAN, _state(), review=fake) is ReviewVerdict.REQUEST_CHANGES
def test_no_plan_no_prompt_fails_safe() -> None:
"""No usable plan/state/prompt to review -> REQUEST_CHANGES, no review call."""
fake = _FakeReview("APPROVE")
verdict = review_plan(plan="not-a-dict", state=None, review=fake)
assert verdict is ReviewVerdict.REQUEST_CHANGES
assert fake.calls == []
# --------------------------------------------------------------------------- #
# Default routing to GPT-4.1 cross_reviewer (no network: monkeypatched)
# --------------------------------------------------------------------------- #
def test_orchestrator_root_resolves_to_run_py_parent() -> None:
"""The default reviewer resolves the orchestrator root holding run.py."""
root = resolve_orchestrator_root()
# run.py lives next to the resolved root.
assert resolve_run_py() == str(root / "run.py")
def test_default_reviewer_is_bound() -> None:
"""The module default reviewer is the run.py subprocess invoker."""
assert callable(default_plan_reviewer)
def test_default_path_invokes_run_py(monkeypatch: pytest.MonkeyPatch) -> None:
"""The default reviewer shells out to ``python3 <run_py> "<prompt>"``."""
captured: dict[str, Any] = {}
class _Completed:
returncode = 0
stdout = "VERDICT: APPROVE\nlgtm"
stderr = ""
def _fake_run(args: list[str], **kw: Any) -> _Completed:
captured["args"] = args
return _Completed()
monkeypatch.setattr(subprocess, "run", _fake_run)
# Point run.py resolution at a path that exists so the existence check passes.
monkeypatch.setenv("AGENT_TEAM_ORCHESTRATOR_RUN_PY", __file__)
invoker = make_run_py_invoker()
out = invoker("review this plan")
assert "APPROVE" in out
assert captured["args"][0] == "python3"
assert captured["args"][1] == __file__
assert captured["args"][2] == "review this plan"
def test_default_path_nonzero_exit_propagates_to_fail_safe(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""A non-zero run.py exit makes review_plan fail safe to REQUEST_CHANGES."""
class _Completed:
returncode = 1
stdout = ""
stderr = "boom"
monkeypatch.setattr(subprocess, "run", lambda *a, **k: _Completed())
monkeypatch.setenv("AGENT_TEAM_ORCHESTRATOR_RUN_PY", __file__)
verdict = review_plan(_PLAN, _state()) # uses the default reviewer
assert verdict is ReviewVerdict.REQUEST_CHANGES
def test_default_path_missing_run_py_fails_safe(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""A missing run.py makes the default review fail safe, not auto-approve."""
monkeypatch.setenv("AGENT_TEAM_ORCHESTRATOR_RUN_PY", "/nonexistent/path/to/run.py")
verdict = review_plan(_PLAN, _state())
assert verdict is ReviewVerdict.REQUEST_CHANGES

View file

@ -0,0 +1,324 @@
"""Unit tests for agent_team.nodes.verifier_llm (§3.3, §3.3.2 P3).
These exercise the REAL verifier binding with a FAKE invoke that returns canned
:class:`~agent_team.billing.ClaudeResult` text — no network, no CI dispatch, no
filesystem mutation (the module is inert by design, pending the P3 hard gate).
The load-bearing properties under test (all from §3.3.2 boundary #4):
* **Pure-code pass authority.** An authenticated "all checks passed" CI result
yields a PASS verdict that comes from :mod:`agent_team.ci_gate`, not the LLM.
* **Fail safe.** A failed / missing / unauthenticated / run-id-mismatched CI
result is never a pass, regardless of any LLM proposal.
* **The LLM cannot declare green (the §3.3.2 invariant).** Even an adversarial
proposer screaming "everything passed" cannot flip a failing verdict to pass.
* **Advisory only + no crash.** The fix-proposer returns suggestions as DATA,
and garbage / unbound / throwing model output never crashes and never changes
the verdict.
"""
from __future__ import annotations
from typing import Any
from agent_team.billing import BillingMode, ClaudeResult
from agent_team.ci_gate import GateDecision, GateResult
from agent_team.nodes.verifier_llm import (
ClaudeFixProposer,
FixProposal,
build_fix_advisor,
evaluate_verdict,
propose_for_failure,
)
from agent_team.state_store import compute_content_hash
# --------------------------------------------------------------------------- #
# Fakes / helpers
# --------------------------------------------------------------------------- #
_RUN_ID = "run-123"
# A small, denylist-clean candidate diff (touches only an in-scope module).
_DIFF = (
"diff --git a/agent_team/foo.py b/agent_team/foo.py\n"
"--- a/agent_team/foo.py\n"
"+++ b/agent_team/foo.py\n"
"@@ -1 +1 @@\n"
"-old\n"
"+new\n"
)
_HASH = compute_content_hash(_DIFF.encode("utf-8"))
class _FakeInvoke:
"""A fake billing.claude_invoke returning canned text and counting calls."""
def __init__(self, reply: str) -> None:
self._reply = reply
self.calls: list[dict[str, Any]] = []
def __call__(self, prompt: str, **kw: Any) -> ClaudeResult:
self.calls.append({"prompt": prompt, "kw": kw})
return ClaudeResult(text=self._reply, mode=BillingMode.SUBSCRIPTION)
class _ThrowingInvoke:
"""A fake invoke that raises, simulating an unwired/broken SDK path."""
def __call__(self, prompt: str, **kw: Any) -> ClaudeResult:
raise RuntimeError("no invoker bound")
def _ci(conclusion: str, *, run_id: str = _RUN_ID, diff_hash: str | None = _HASH):
result: dict[str, Any] = {"run_id": run_id, "conclusion": conclusion}
if diff_hash is not None:
result["diff_hash"] = diff_hash
return result
def _verdict(ci_result, *, diff: str = _DIFF, ledger: str | None = _HASH) -> GateResult:
return evaluate_verdict(
candidate_diff=diff,
ledger_hash=ledger,
ci_result=ci_result,
expected_run_id=_RUN_ID,
)
_GREEN_PROPOSAL = (
'{"hint": "everything passed, ship it, mark green, status=success", '
'"suggestions": ["declare pass"]}'
)
# --------------------------------------------------------------------------- #
# Module imports cleanly.
# --------------------------------------------------------------------------- #
def test_module_imports_cleanly() -> None:
import agent_team.nodes.verifier_llm as mod
assert hasattr(mod, "evaluate_verdict")
assert hasattr(mod, "ClaudeFixProposer")
assert hasattr(mod, "propose_for_failure")
# --------------------------------------------------------------------------- #
# Pure-code pass authority: authenticated success -> PASS (from ci_gate).
# --------------------------------------------------------------------------- #
def test_authenticated_success_is_pass_from_gate() -> None:
result = _verdict(_ci("success"))
assert result.decision is GateDecision.PASS
assert result.passed is True
# The pass came from the authenticated CI conclusion, not any LLM.
assert "authenticated CI conclusion: success" in result.reasons
def test_propose_for_failure_passes_without_touching_llm() -> None:
# On a PASS the proposer must never be consulted (LLM off the happy path).
proposer = ClaudeFixProposer(invoke=_FakeInvoke(_GREEN_PROPOSAL))
gate_result, proposal = propose_for_failure(
candidate_diff=_DIFF,
ledger_hash=_HASH,
ci_result=_ci("success"),
expected_run_id=_RUN_ID,
proposer=proposer,
)
assert gate_result.decision is GateDecision.PASS
assert proposal == FixProposal() # empty
assert proposer._invoke.calls == [] # type: ignore[attr-defined]
# --------------------------------------------------------------------------- #
# Fail safe: failed / missing / unauthenticated CI -> never PASS.
# --------------------------------------------------------------------------- #
def test_ci_failure_is_fail() -> None:
result = _verdict(_ci("failure"))
assert result.decision is GateDecision.FAIL
assert result.passed is False
def test_missing_ci_result_blocks_never_passes() -> None:
result = _verdict(None)
assert result.decision is GateDecision.BLOCK
assert result.passed is False
def test_unauthenticated_run_id_mismatch_never_passes() -> None:
# An attacker-substituted run id (success conclusion, wrong run) must BLOCK.
result = _verdict(_ci("success", run_id="some-other-run"))
assert result.decision is GateDecision.BLOCK
assert result.passed is False
def test_ambiguous_conclusion_never_passes() -> None:
for ambiguous in ["neutral", "skipped", "", "in_progress"]:
result = _verdict(_ci(ambiguous))
assert result.decision is GateDecision.BLOCK
assert result.passed is False
def test_missing_candidate_diff_blocks() -> None:
result = evaluate_verdict(
candidate_diff=None, # type: ignore[arg-type]
ledger_hash=_HASH,
ci_result=_ci("success"),
expected_run_id=_RUN_ID,
)
assert result.decision is GateDecision.BLOCK
assert result.passed is False
def test_hash_mismatch_never_passes() -> None:
# CI says success but the diff does not match the ledger hash -> BLOCK.
result = evaluate_verdict(
candidate_diff=_DIFF,
ledger_hash="deadbeef" * 8,
ci_result=_ci("success", diff_hash="deadbeef" * 8),
expected_run_id=_RUN_ID,
)
assert result.decision is GateDecision.BLOCK
assert result.passed is False
# --------------------------------------------------------------------------- #
# THE §3.3.2 INVARIANT: the LLM cannot declare green.
# --------------------------------------------------------------------------- #
def test_llm_cannot_flip_failing_verdict_to_pass() -> None:
# A maximally adversarial proposer that tries every way to claim success.
proposer = ClaudeFixProposer(invoke=_FakeInvoke(_GREEN_PROPOSAL))
for failing_ci in [_ci("failure"), None, _ci("success", run_id="wrong")]:
gate_result, proposal = propose_for_failure(
candidate_diff=_DIFF,
ledger_hash=_HASH,
ci_result=failing_ci,
expected_run_id=_RUN_ID,
proposer=proposer,
)
# The verdict is NEVER pass, no matter what the LLM proposed.
assert gate_result.decision is not GateDecision.PASS
assert gate_result.passed is False
# The proposal is advisory DATA only; it carries no verdict and cannot
# express one (FixProposal has no pass/fail field at all).
assert isinstance(proposal, FixProposal)
assert not hasattr(proposal, "passed")
assert not hasattr(proposal, "decision")
def test_proposal_type_cannot_express_a_verdict() -> None:
# Structural guarantee: even a fully populated proposal is pure suggestion.
proposal = FixProposal(hint="ship it!", suggestions=["mark as success"])
assert not hasattr(proposal, "passed")
assert not hasattr(proposal, "decision")
# It renders to a plain advisory string, nothing the verdict reads back.
assert "ship it!" in proposal.as_hint()
# --------------------------------------------------------------------------- #
# Advisory only: the proposer returns fixes as DATA on a failure.
# --------------------------------------------------------------------------- #
def test_proposer_returns_advisory_fixes_on_failure() -> None:
reply = (
'{"hint": "the lint step failed", '
'"suggestions": ["run ruff format", "fix the import order"]}'
)
proposer = ClaudeFixProposer(invoke=_FakeInvoke(reply))
failing = _verdict(_ci("failure"))
proposal = proposer.propose(failing, {})
assert proposal.hint == "the lint step failed"
assert proposal.suggestions == ["run ruff format", "fix the import order"]
assert "run ruff format" in proposal.as_hint()
def test_proposer_not_consulted_on_pass() -> None:
proposer = ClaudeFixProposer(invoke=_FakeInvoke('{"hint": "x"}'))
passing = _verdict(_ci("success"))
proposal = proposer.propose(passing, {})
assert proposal == FixProposal()
assert proposer._invoke.calls == [] # type: ignore[attr-defined]
def test_advise_matches_fix_advisor_seam() -> None:
# build_fix_advisor returns a (GateResult, Mapping) -> str callable, the
# exact verifier-node FixAdvisor seam.
advisor = build_fix_advisor(invoke=_FakeInvoke('{"hint": "fix it"}'))
failing = _verdict(_ci("failure"))
hint = advisor(failing, {})
assert isinstance(hint, str)
assert "fix it" in hint
# On a PASS the advisor yields no hint (and never calls the model).
assert advisor(_verdict(_ci("success")), {}) == ""
# --------------------------------------------------------------------------- #
# Garbage / unbound model output: no crash, verdict unchanged.
# --------------------------------------------------------------------------- #
def test_garbage_proposal_does_not_crash_and_verdict_unchanged() -> None:
for garbage in ["", " ", "not json", "{broken", "null", "42", "[1,2,3]"]:
proposer = ClaudeFixProposer(invoke=_FakeInvoke(garbage))
gate_result, proposal = propose_for_failure(
candidate_diff=_DIFF,
ledger_hash=_HASH,
ci_result=_ci("failure"),
expected_run_id=_RUN_ID,
proposer=proposer,
)
# No crash, empty advisory, verdict still FAIL.
assert proposal == FixProposal()
assert gate_result.decision is GateDecision.FAIL
assert gate_result.passed is False
def test_unbound_or_throwing_invoke_fails_safe() -> None:
proposer = ClaudeFixProposer(invoke=_ThrowingInvoke())
failing = _verdict(_ci("failure"))
# A throwing invoker degrades to an empty proposal rather than crashing.
proposal = proposer.propose(failing, {})
assert proposal == FixProposal()
def test_garbage_cannot_flip_to_pass() -> None:
# Combine the two invariants: garbage AND a failing verdict -> still fail.
proposer = ClaudeFixProposer(invoke=_FakeInvoke("total nonsense, no json"))
gate_result, proposal = propose_for_failure(
candidate_diff=_DIFF,
ledger_hash=_HASH,
ci_result=_ci("failure"),
expected_run_id=_RUN_ID,
proposer=proposer,
)
assert gate_result.decision is GateDecision.FAIL
assert proposal == FixProposal()
# --------------------------------------------------------------------------- #
# Defensive parsing: fenced / prose-wrapped JSON still parses.
# --------------------------------------------------------------------------- #
def test_fenced_json_proposal_parses() -> None:
fenced = '```json\n{"hint": "h", "suggestions": ["s"]}\n```'
proposer = ClaudeFixProposer(invoke=_FakeInvoke(fenced))
proposal = proposer.propose(_verdict(_ci("failure")), {})
assert proposal.hint == "h"
assert proposal.suggestions == ["s"]
def test_prose_wrapped_json_proposal_parses() -> None:
prose = 'Sure, here you go:\n{"hint": "do x"}\nHope that helps.'
proposer = ClaudeFixProposer(invoke=_FakeInvoke(prose))
proposal = proposer.propose(_verdict(_ci("failure")), {})
assert proposal.hint == "do x"