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.
292 lines
13 KiB
Python
292 lines
13 KiB
Python
"""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)
|