Merge pull request #13 from Sea-Haven-Industries/feature/agent-team-plane2-p3-p4

agent-team Plane-2: P4 live transports + intake, P3-inert build/verify subgraph
This commit is contained in:
Adam Moussa 2026-06-18 13:40:48 -04:00 • committed by GitHub
commit 73e35f3acd
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
13 changed files with 2892 additions and 21 deletions

View file

@ -66,6 +66,7 @@ if TYPE_CHECKING: # pragma: no cover - typing only
__all__ = [ __all__ = [
"Coordinator", "Coordinator",
"build_verify_wiring",
"default_clarify_node_factory", "default_clarify_node_factory",
] ]
@ -96,6 +97,19 @@ ReviewWiring = Callable[
[], [],
"tuple[Callable[[PipelineState], PipelineState], Callable[[PipelineState], str]]", "tuple[Callable[[PipelineState], PipelineState], Callable[[PipelineState], str]]",
] ]
# Composes the opt-in P3 build->verify subgraph and yields the
# ``(build_node, verify_node, route_after_verify)`` tuple
# :func:`agent_team.graph.build_graph` wires off the review loop's "build" route.
# None -> no P3 subgraph (P2: the review's "build" route is the terminus). This
# stays OPT-IN and INERT: the production run-team path never injects it (it is
# held until the §3.3.2 CI trust-boundary security gate clears).
BuildVerifyWiring = Callable[
[],
"tuple["
"Callable[[PipelineState], dict[str, Any]], "
"Callable[[PipelineState], PipelineState], "
"Callable[[PipelineState], str]]",
]
# A checkpointer factory over the db path: returns the BaseCheckpointSaver the # A checkpointer factory over the db path: returns the BaseCheckpointSaver the
# graph is compiled with. Defaults to the production SQLite checkpointer # graph is compiled with. Defaults to the production SQLite checkpointer
@ -189,6 +203,57 @@ def default_review_wiring() -> tuple[
return review_loop.bind_review_node(), review_loop.route_after_review return review_loop.bind_review_node(), review_loop.route_after_review
def build_verify_wiring() -> tuple[
Callable[[PipelineState], "dict[str, Any]"],
Callable[[PipelineState], PipelineState],
Callable[[PipelineState], str],
]:
"""Compose the OPT-IN, INERT P3 build->verify subgraph (held for the gate).
Returns the ``(build_node, verify_node, route_after_verify)`` tuple
:func:`agent_team.graph.build_graph` hangs off the review loop's "build"
route, composed from :mod:`agent_team.nodes.build_verify_subgraph` with its
INERT defaults:
* the BUILD node is built with NO injected diff builder
(:func:`~agent_team.nodes.build_verify_subgraph.make_build_node` with
``diff_builder=None``), so it falls back to the committed default builder
that fails loudly in an un-wired environment rather than emitting an empty
diff; and
* the VERIFY node is built with NO ``ci_result`` fetcher
(:func:`~agent_team.nodes.build_verify_subgraph.make_verify_node` with the
default ``_no_ci_result`` -> ``None``), so the pure-code gate has no
authenticated pass to read, returns ``BLOCK``, and the task parks.
This is the FAIL-SAFE composition: with no authenticated CI result the
pipeline can NEVER fabricate a pass, so wiring this subgraph in is inert
until a leaf binds the real diff-builder + read-only-PAT CI fetcher via
:func:`~agent_team.nodes.build_verify_subgraph.bind_diff_builder` /
:func:`~agent_team.nodes.build_verify_subgraph.bind_ci_result_fetcher` AFTER
the §3.3.2 CI trust-boundary clears ``/sh-security-review`` + the GPT-4.1
cross-review. The production ``run-team.py`` path deliberately does NOT inject
this factory; it stays opt-in/off.
Lazy-imported for the same import-hygiene reason as the clarifier / planner /
review factories (the subgraph pulls in the builders + verifier leaves).
"""
from agent_team.nodes.build_verify_subgraph import (
make_build_node,
make_verify_node,
route_after_verify,
)
from agent_team.nodes.verifier import VerifierConfig
# INERT VerifierConfig: ``expected_run_id`` is required by the dataclass but
# is moot while the default fetcher yields no authenticated CI result (the
# gate BLOCKs and the task parks regardless). A real leaf supplies the task's
# actual expected run id once the gate clears.
config = VerifierConfig(expected_run_id="")
build_node = make_build_node(diff_builder=None)
verify_node = make_verify_node(config, ci_result_fetcher=None)
return build_node, verify_node, route_after_verify
class Coordinator: class Coordinator:
"""Owns the live Plane-2 runtime: graph + resume worker + transport (§3.3). """Owns the live Plane-2 runtime: graph + resume worker + transport (§3.3).
@ -211,6 +276,7 @@ class Coordinator:
build_clarify_node: ClarifyNodeFactory | None = None, build_clarify_node: ClarifyNodeFactory | None = None,
build_plan_node: PlanNodeFactory | None = None, build_plan_node: PlanNodeFactory | None = None,
review_wiring: ReviewWiring | None = None, review_wiring: ReviewWiring | None = None,
build_verify_wiring: BuildVerifyWiring | None = None,
build_checkpointer: CheckpointerFactory | None = None, build_checkpointer: CheckpointerFactory | None = None,
resume_queue: "queue.Queue[Any] | None" = None, resume_queue: "queue.Queue[Any] | None" = None,
deadline_window: timedelta | None = None, deadline_window: timedelta | None = None,
@ -224,6 +290,10 @@ class Coordinator:
# path injects default_plan_node_factory + default_review_wiring. # path injects default_plan_node_factory + default_review_wiring.
self._build_plan_node = build_plan_node self._build_plan_node = build_plan_node
self._review_wiring = review_wiring self._review_wiring = review_wiring
# P3 (build->verify) is OPT-IN and INERT: left None, the graph stops at
# the P2 review terminus. The production run-team path never injects it;
# it is held until the §3.3.2 CI trust-boundary security gate clears.
self._build_verify_wiring = build_verify_wiring
self._build_checkpointer = ( self._build_checkpointer = (
build_checkpointer or graph_mod.build_sqlite_checkpointer build_checkpointer or graph_mod.build_sqlite_checkpointer
) )
@ -294,12 +364,20 @@ class Coordinator:
if self._review_wiring is not None: if self._review_wiring is not None:
review_node, route_review = self._review_wiring() review_node, route_review = self._review_wiring()
# P3 (opt-in/INERT): the build->verify subgraph tuple, only wired when
# the review loop is also wired (it hangs off the review's "build"
# route). None in the production path.
build_verify: Any = None
if self._build_verify_wiring is not None:
build_verify = self._build_verify_wiring()
self._graph = graph_mod.build_graph( self._graph = graph_mod.build_graph(
checkpointer, checkpointer,
live_clarify_node=clarify_node, live_clarify_node=clarify_node,
live_plan_node=plan_node, live_plan_node=plan_node,
review_node=review_node, review_node=review_node,
route_review=route_review, route_review=route_review,
build_verify=build_verify,
) )
# The ResumeWorker is satisfied directly by the compiled LangGraph app # The ResumeWorker is satisfied directly by the compiled LangGraph app
@ -354,7 +432,14 @@ class Coordinator:
if self._graph is None: if self._graph is None:
raise RuntimeError("Coordinator.start_task called before setup()") raise RuntimeError("Coordinator.start_task called before setup()")
_LOG.info("start_task intake (transport=%s): %s", transport_name, task_text) # Intake text is untrusted (a GitHub issue body, etc.); truncate and
# escape newlines so a forged multi-line body cannot spoof operator log
# lines (log-injection hygiene).
_LOG.info(
"start_task intake (transport=%s): %s",
transport_name,
task_text[:200].replace("\n", "\\n").replace("\r", "\\r"),
)
thread_id, _state = graph_mod.start_task(self._graph, transport=transport_name) thread_id, _state = graph_mod.start_task(self._graph, transport=transport_name)
question = graph_mod.pending_question(self._graph, thread_id=thread_id) question = graph_mod.pending_question(self._graph, thread_id=thread_id)

View file

@ -62,6 +62,8 @@ if TYPE_CHECKING: # pragma: no cover - typing only
from langgraph.graph.state import CompiledStateGraph from langgraph.graph.state import CompiledStateGraph
__all__ = [ __all__ = [
"APPROVED_ROUTE",
"BUILD_NODE",
"BUILD_ROUTE", "BUILD_ROUTE",
"CLARIFY", "CLARIFY",
"DEFAULT_CLARIFY_DEADLINE", "DEFAULT_CLARIFY_DEADLINE",
@ -70,6 +72,7 @@ __all__ = [
"PARKED_ROUTE", "PARKED_ROUTE",
"PLAN", "PLAN",
"REVIEW", "REVIEW",
"VERIFY_NODE",
"build_graph", "build_graph",
"build_sqlite_checkpointer", "build_sqlite_checkpointer",
"clarify_node", "clarify_node",
@ -99,6 +102,25 @@ REVIEW = "review"
BUILD_ROUTE = "build" BUILD_ROUTE = "build"
PARKED_ROUTE = "parked" PARKED_ROUTE = "parked"
# P3 (build -> verify subgraph) vertex ids. These are the GRAPH VERTEX names the
# opt-in P3 subgraph hangs off the review loop's "build" route; they are kept
# distinct from the route-id constants above (BUILD_ROUTE / PARKED_ROUTE) and
# from PLAN/REVIEW so the conditional-edge maps never collide a route key with a
# vertex id. The subgraph itself is supplied wholesale by the Integrate-phase
# caller (the ``build_verify`` tuple), so graph.py does not import
# agent_team.nodes.build_verify_subgraph (no wiring import cycle); it only owns
# the topology that connects the injected nodes.
BUILD_NODE = "build_node"
VERIFY_NODE = "verify_node"
# P3 route ids returned by the injected ``route_after_verify`` function. They
# mirror agent_team.nodes.build_verify_subgraph.APPROVED_ROUTE / BUILD_ROUTE /
# PARKED_ROUTE by VALUE so this module wires the VERIFY conditional-edge map
# without importing that module. APPROVED_ROUTE is the build->verify-specific
# PASS terminus (the draft-PR endpoint); BUILD_ROUTE loops back to the builders
# under the build-loop budget; PARKED_ROUTE is the fail-safe escalation.
APPROVED_ROUTE = "approved"
# The P1 stage order (§7.1): intake -> clarify -> plan, then stop. Builders and # The P1 stage order (§7.1): intake -> clarify -> plan, then stop. Builders and
# verifiers (BUILD/VERIFY) are deliberately NOT wired here — P1 ends at an # verifiers (BUILD/VERIFY) are deliberately NOT wired here — P1 ends at an
# approved plan with no build (§7.1 "Stops at an approved plan, no build yet"). # approved plan with no build (§7.1 "Stops at an approved plan, no build yet").
@ -266,6 +288,12 @@ def build_graph(
live_plan_node: Callable[[PipelineState], PipelineState] | None = None, live_plan_node: Callable[[PipelineState], PipelineState] | None = None,
review_node: Callable[[PipelineState], PipelineState] | None = None, review_node: Callable[[PipelineState], PipelineState] | None = None,
route_review: Callable[[PipelineState], str] | None = None, route_review: Callable[[PipelineState], str] | None = None,
build_verify: tuple[
Callable[[PipelineState], dict[str, Any]],
Callable[[PipelineState], PipelineState],
Callable[[PipelineState], str],
]
| None = None,
) -> CompiledStateGraph: ) -> CompiledStateGraph:
"""Assemble + compile the P1 pipeline ``StateGraph`` (§3.3, §7.1). """Assemble + compile the P1 pipeline ``StateGraph`` (§3.3, §7.1).
@ -311,6 +339,27 @@ def build_graph(
``review_node`` requires ``route_review`` (and a real ``live_plan_node`` that ``review_node`` requires ``route_review`` (and a real ``live_plan_node`` that
advances to REVIEW); passing one without the other is a wiring error. advances to REVIEW); passing one without the other is a wiring error.
``build_verify`` wires the **OPT-IN P3** build -> verify subgraph and is
supplied wholesale as the tuple
``(build_node, verify_node, route_after_verify)`` the coordinator composes
from :mod:`agent_team.nodes.build_verify_subgraph` (passed in so this module
never imports that module — no wiring import cycle). It only takes effect
when the review loop is also wired (it hangs off the review's ``"build"``
route):
* **P2 (default):** ``build_verify`` is ``None`` -> the review's ``"build"``
route terminates at ``END`` (the approved-plan terminus), exactly as
before. Production stays P2 (clarify -> plan -> review).
* **P3:** ``build_verify`` is given (with ``review_node``) -> the review's
``"build"`` route is REPOINTED at the BUILD node, ``BUILD -> VERIFY`` is
wired, and ``route_after_verify`` maps ``{approved -> END (PR terminus),
build -> BUILD (bounded build<->verify loop), parked -> END (escalation)}``.
The subgraph stays INERT unless the caller binds real diff-builder / CI
seams (held for the §3.3.2 security gate); with the default INERT seams the
verifier gate has no authenticated pass and parks. Passing
``build_verify`` without ``review_node`` is a wiring error (there is no
``"build"`` route to repoint).
""" """
clarify = live_clarify_node if live_clarify_node is not None else clarify_node clarify = live_clarify_node if live_clarify_node is not None else clarify_node
plan = live_plan_node if live_plan_node is not None else plan_node plan = live_plan_node if live_plan_node is not None else plan_node
@ -321,6 +370,13 @@ def build_graph(
"function, e.g. review_loop.route_after_review)." "function, e.g. review_loop.route_after_review)."
) )
if build_verify is not None and review_node is None:
raise ValueError(
"build_graph: build_verify (the P3 build->verify subgraph) requires "
"review_node — it hangs off the review loop's 'build' route, so there "
"is nothing to repoint without a review loop."
)
builder: StateGraph = StateGraph(PipelineState) builder: StateGraph = StateGraph(PipelineState)
builder.add_node(INTAKE, intake_node) builder.add_node(INTAKE, intake_node)
builder.add_node(CLARIFY, clarify) builder.add_node(CLARIFY, clarify)
@ -334,14 +390,38 @@ def build_graph(
# P1: the plan stage is the terminus. # P1: the plan stage is the terminus.
builder.add_edge(PLAN, END) builder.add_edge(PLAN, END)
else: else:
# P2: plan -> review -> {loop-back to plan | END}. # P2/P3: plan -> review -> {loop-back to plan | build | END}.
builder.add_node(REVIEW, review_node) builder.add_node(REVIEW, review_node)
builder.add_edge(PLAN, REVIEW) builder.add_edge(PLAN, REVIEW)
builder.add_conditional_edges(
REVIEW, if build_verify is None:
route_review, # P2: the review's "build" route is the approved-plan terminus.
{BUILD_ROUTE: END, PLAN: PLAN, PARKED_ROUTE: END}, builder.add_conditional_edges(
) REVIEW,
route_review,
{BUILD_ROUTE: END, PLAN: PLAN, PARKED_ROUTE: END},
)
else:
# P3 (opt-in): repoint the review's "build" route at the BUILD node,
# wire BUILD -> VERIFY, and route the verifier verdict to
# {approved -> END (PR terminus), build -> BUILD (loop), parked ->
# END (escalation)}. The subgraph nodes + router are injected (the
# ``build_verify`` tuple) so this module imports no P3 code.
build_node, verify_node, route_after_verify = build_verify
builder.add_node(BUILD_NODE, build_node)
builder.add_node(VERIFY_NODE, verify_node)
builder.add_conditional_edges(
REVIEW,
route_review,
{BUILD_ROUTE: BUILD_NODE, PLAN: PLAN, PARKED_ROUTE: END},
)
builder.add_edge(BUILD_NODE, VERIFY_NODE)
builder.add_conditional_edges(
VERIFY_NODE,
route_after_verify,
{APPROVED_ROUTE: END, BUILD_ROUTE: BUILD_NODE, PARKED_ROUTE: END},
)
if checkpointer is None: if checkpointer is None:
return builder.compile() return builder.compile()

View file

@ -0,0 +1,286 @@
"""P3-INERT build -> verify subgraph TOPOLOGY (design §3.3, §7.1 P3).
This module is the **wiring topology** for the Plane-2 build -> verify stage::
... -> REVIEW (route "build") -> BUILD -> VERIFY -> {approved | build | parked}
It produces the BUILD node, the VERIFY node, and the
:func:`route_after_verify` conditional-edge function so the Integrate phase can
hang them off :func:`agent_team.graph.build_graph` as a subgraph reachable from
the review loop's ``"build"`` route. It assembles NOTHING by itself: it does not
call :func:`agent_team.graph.build_graph`, and the production default pipeline
stays P2 (clarify -> plan -> review). Hooking this subgraph in is a deliberate,
opt-in Integrate-phase edit.
============================== INERT / HARD-GATE ==========================
P3 (builders + verifier) is HARD-GATED behind ``/sh-security-review`` + a
GPT-4.1 cross-review of the §3.3.2 CI apply/verify trust boundary BEFORE it goes
live. This module is TOPOLOGY + SEAMS ONLY and MUST stay INERT:
* **No live CI.** The VERIFY node consumes an INJECTED ``ci_result`` seam — a
fetcher callable that, given the task state, returns the authenticated CI
conclusion as DATA (exactly what :func:`agent_team.ci_gate.evaluate_ci_gate`
expects). The DEFAULT fetcher returns ``None`` (the current pre-live-CI
reality). It performs NO live CI dispatch, NO OIDC, NO network to GitHub
Actions, NO ``git``/patch apply, and NO filesystem mutation.
* **Fail-safe verdict.** With no ``ci_result`` (the default), the pure-code
gate (:mod:`agent_team.ci_gate`) returns ``BLOCK`` — there is no
authenticated pass to be had — and :func:`route_after_verify` routes the
task to PARKED. The pipeline NEVER fabricates a pass; the gate is the sole
pass authority.
* **LLM stays a fix-proposer.** The verifier's LLM seam
(:data:`agent_team.nodes.verifier.FixAdvisor`) is consulted ONLY on a
failure to author a fix hint. It is structurally incapable of flipping the
verdict to pass (the verdict is computed first, by the gate, and is never
read back from the proposer — see :mod:`agent_team.nodes.verifier_llm`).
The live apply/verify path (the gated wiring of a real diff builder + a real CI
result fetcher) is held for the separate security-review + cross-review gate and
is NOT shipped or enabled here. :func:`bind_diff_builder` and
:func:`bind_ci_result_fetcher` are the injection points a leaf will use to bind
those real seams once the gate clears.
============================================================================
What this module owns (topology + seams only):
* :data:`APPROVED_ROUTE` / :data:`BUILD_ROUTE` / :data:`PARKED_ROUTE` — the
route ids :func:`route_after_verify` returns. ``BUILD_ROUTE`` /
``PARKED_ROUTE`` mirror :data:`agent_team.graph.BUILD_ROUTE` /
:data:`agent_team.graph.PARKED_ROUTE` by VALUE so the conditional-edge map the
Integrate phase builds matches without this module importing ``graph`` (which
would be a wiring import cycle).
* :func:`make_build_node` — factory producing the single-argument BUILD node,
threading an injectable :class:`~agent_team.nodes.builders.DiffBuilder` into
:func:`agent_team.nodes.builders.builders_node`.
* :func:`make_verify_node` — factory producing the single-argument VERIFY node,
threading an injectable ``ci_result`` fetcher into
:func:`agent_team.nodes.verifier.verifier_node` (default fetcher -> ``None``).
* :func:`route_after_verify` — the LangGraph conditional-edge function that
reads the verdict the VERIFY node recorded and returns the next route id.
* :func:`bind_diff_builder` / :func:`bind_ci_result_fetcher` — the gated-live
injection points (held for the security gate).
It imports the committed node + foundation contracts verbatim and redefines none
of them. No SDK is imported at module top (deferred discipline mirroring
:func:`agent_team.graph.build_sqlite_checkpointer`); it is fully unit-testable
with no network.
"""
from __future__ import annotations
from collections.abc import Callable, Mapping
from typing import Any
from agent_team.nodes.builders import DiffBuilder, builders_node
from agent_team.nodes.verifier import VerifierConfig, verifier_node
from agent_team.task_model import Phase, PipelineState
__all__ = [
"APPROVED_ROUTE",
"BUILD_NODE",
"BUILD_ROUTE",
"CiResultFetcher",
"PARKED_ROUTE",
"VERIFY_NODE",
"bind_ci_result_fetcher",
"bind_diff_builder",
"make_build_node",
"make_verify_node",
"route_after_verify",
]
# --- Node names (graph vertices). ------------------------------------------
# Kept as constants so the Integrate-phase wiring references the subgraph
# vertices by name rather than by string literal.
BUILD_NODE = "build"
VERIFY_NODE = "verify"
# --- Route ids returned by route_after_verify. -----------------------------
# These mirror agent_team.graph.BUILD_ROUTE / PARKED_ROUTE by VALUE so the
# conditional-edge map the Integrate phase builds lines up without importing
# graph here (that would be a wiring import cycle). APPROVED_ROUTE is the
# build->verify-specific PASS terminus (the draft-PR endpoint); BUILD_ROUTE is
# the loop-back to the builders on a recoverable failure; PARKED_ROUTE is the
# fail-safe escalation (the only reachable route while INERT, since the default
# fetcher yields no authenticated pass).
APPROVED_ROUTE = "approved"
BUILD_ROUTE = "build"
PARKED_ROUTE = "parked"
# The injectable CI-result seam: given the task state, return the authenticated,
# patch-independent CI conclusion as a mapping (run_id / conclusion / diff_hash),
# or ``None`` when there is no authenticated result. The DEFAULT
# (:func:`_no_ci_result`) always returns ``None`` (the INERT pre-live-CI
# reality), so the gate BLOCKs and the task parks — never a fabricated pass. The
# real fetcher (read-only PAT against the GitHub Checks/Actions API) is bound via
# :func:`bind_ci_result_fetcher` only after the §3.3.2 trust boundary clears its
# security gate.
CiResultFetcher = Callable[[PipelineState], Mapping[str, Any] | None]
def _no_ci_result(state: PipelineState) -> None:
"""Default :data:`CiResultFetcher`: there is NO authenticated CI result.
This is the INERT, pre-live-CI reality. Returning ``None`` means the
pure-code gate (:func:`agent_team.ci_gate.evaluate_ci_gate`) has no
authenticated conclusion to read and therefore returns ``BLOCK`` — never a
pass. The subgraph thus fails SAFE to PARKED until a real fetcher is bound
via :func:`bind_ci_result_fetcher` (which is held for the security gate).
"""
return None
def make_build_node(
*,
diff_builder: DiffBuilder | None = None,
config: Mapping[str, Any] | None = None,
) -> Callable[[PipelineState], dict[str, Any]]:
"""Produce the single-argument BUILD node (approved plan -> candidate diff).
Wraps :func:`agent_team.nodes.builders.builders_node` as a one-argument
``PipelineState -> partial PipelineState`` closure so LangGraph can add it as
a vertex without seeing the node's ``builder`` / ``config`` keyword params
(LangGraph would otherwise try to inject its own ``RunnableConfig`` there —
the same hazard :func:`agent_team.nodes.review_loop.bind_review_node`
guards against). The injected ``diff_builder`` is threaded straight to the
node's :class:`~agent_team.nodes.builders.DiffBuilder` seam.
INERT: when ``diff_builder`` is ``None`` the node falls back to its committed
default (:func:`agent_team.nodes.builders.default_diff_builder`), which fails
LOUDLY in an un-wired environment (the billing seam raises until configured)
rather than emitting an empty diff. The real DeepSeek path is bound via
:func:`bind_diff_builder` once the §3.3.2 gate clears. The node itself never
applies a patch — it emits the diff as DATA plus the box-side
trust-control-surface scan + integrity hash.
"""
def node(state: PipelineState) -> dict[str, Any]:
return builders_node(state, builder=diff_builder, config=config)
return node
def make_verify_node(
config: VerifierConfig,
*,
ci_result_fetcher: CiResultFetcher | None = None,
) -> Callable[[PipelineState], PipelineState]:
"""Produce the single-argument VERIFY node (gate the CI result, decide phase).
Wraps :func:`agent_team.nodes.verifier.verifier_node` as a one-argument
closure over ``config`` (a :class:`~agent_team.nodes.verifier.VerifierConfig`)
so it wires straight into LangGraph without a manage-injected ``config``
param. The INERT ``ci_result`` seam is the key here: the node reads its CI
conclusion from ``state["ci_results"]``, so this wrapper FETCHES that result
via the injected ``ci_result_fetcher`` and merges it into the state BEFORE
delegating to the node.
The default fetcher (:func:`_no_ci_result`) returns ``None`` (the pre-live-CI
reality). With no authenticated CI result the pure-code gate returns
``BLOCK`` and the node parks the task — it can NEVER fabricate a pass. The
LLM verifier seam stays a fix-PROPOSER only (the gate is the sole pass
authority); see :mod:`agent_team.nodes.verifier_llm`.
The fetcher is called defensively: it receives the task state and returns the
authenticated CI conclusion mapping (``run_id`` / ``conclusion`` /
``diff_hash``) or ``None``. Any value other than a mapping is treated as
"no result" (``None``), so a malformed fetcher fails SAFE to BLOCK rather
than smuggling something past the gate. The real read-only-PAT fetcher is
bound via :func:`bind_ci_result_fetcher` only after the §3.3.2 trust boundary
clears its security gate.
"""
fetcher: CiResultFetcher = (
ci_result_fetcher if ci_result_fetcher is not None else _no_ci_result
)
def node(state: PipelineState) -> PipelineState:
fetched = fetcher(state)
ci_result = fetched if isinstance(fetched, Mapping) else None
# Merge the (possibly None) fetched CI result into the state the node
# reads from, WITHOUT mutating the caller's state object. The node reads
# ``ci_results``; a None result leaves the gate with nothing to pass on.
scoped_state: dict[str, Any] = dict(state)
scoped_state["ci_results"] = ci_result
return verifier_node(scoped_state, config)
return node
def route_after_verify(state: PipelineState) -> str:
"""LangGraph conditional-edge: the next route id after the VERIFY node.
Reads the phase the VERIFY node recorded (the pure-code gate's verdict,
already merged into ``current_phase`` / ``status``) and maps it to a route
id the Integrate-phase conditional-edge map keys against:
* gate PASS -> phase ``DONE`` -> :data:`APPROVED_ROUTE` (the draft-PR
terminus). UNREACHABLE while INERT — the default fetcher yields no
authenticated pass, so the gate never returns PASS.
* gate FAIL under the build-loop budget -> phase ``BUILD`` ->
:data:`BUILD_ROUTE` (loop back to the builders with the fix hint).
* gate BLOCK, or FAIL at/over the budget -> phase ``PARKED`` ->
:data:`PARKED_ROUTE` (escalate to human + GPT cross-review; ALARM).
FAILS SAFE: any unexpected / missing phase routes to
:data:`PARKED_ROUTE` rather than advancing, so an ambiguous state parks for a
human instead of shipping. The verdict is owned entirely by the gate (the
node already applied it); this function only reads the recorded phase.
"""
phase = state.get("current_phase")
if phase == Phase.DONE.value:
return APPROVED_ROUTE
if phase == Phase.BUILD.value:
return BUILD_ROUTE
if phase == Phase.PARKED.value:
return PARKED_ROUTE
# Unknown / missing phase (the node always sets one of the above) -> park
# fail-closed rather than advancing an ambiguous state.
return PARKED_ROUTE
def bind_diff_builder(
diff_builder: DiffBuilder,
) -> Callable[[PipelineState], dict[str, Any]]:
"""Bind the REAL diff builder into a BUILD node (GATED-LIVE injection point).
Thin convenience over :func:`make_build_node` for the leaf that, once the
§3.3.2 trust boundary clears ``/sh-security-review`` + the GPT-4.1
cross-review, binds the real DeepSeek ``fast_coder`` path (adapt
:func:`agent_team.nodes.builders_llm.as_diff_builder` into a
:class:`~agent_team.nodes.builders.DiffBuilder`). Binding it does NOT enable
any apply/verify behaviour — the BUILD node still only EMITS a diff as DATA
plus the box-side scan + hash. Held for the gate; not wired here.
"""
return make_build_node(diff_builder=diff_builder)
def bind_ci_result_fetcher(
config: VerifierConfig,
ci_result_fetcher: CiResultFetcher,
) -> Callable[[PipelineState], PipelineState]:
"""Bind the REAL CI-result fetcher into a VERIFY node (GATED-LIVE injection).
Thin convenience over :func:`make_verify_node` for the leaf that, once the
§3.3.2 trust boundary clears its security gate, binds the real authenticated
CI-result fetcher (read-only PAT against the GitHub Checks/Actions API, NOT
enabled here). The fetcher returns the authenticated conclusion as DATA;
pass/fail remains owned by the pure-code gate, so binding a fetcher only
GIVES the gate a result to read — it can never make the LLM the pass
authority. Held for the gate; not wired here.
"""
return make_verify_node(config, ci_result_fetcher=ci_result_fetcher)
# A module-level note for the Integrate phase (no execution): the build->verify
# subgraph is hung off the review loop's "build" route. The conditional-edge map
# from VERIFY should send APPROVED_ROUTE to the PR/draft terminus, BUILD_ROUTE
# back to the BUILD node (the bounded build<->verify loop, capped by
# VerifierConfig.max_build_loops), and PARKED_ROUTE to the escalation terminus.
# build_graph wires this in opt-in; this module never assembles it itself.
_INTEGRATE_NOTE = (
"review('build') -> BUILD -> VERIFY -> route_after_verify -> "
"{approved: PR terminus, build: BUILD (loop), parked: escalation}"
)

View file

@ -0,0 +1,265 @@
"""Live Claude-Code-on-the-Mac delivery wiring (design §3.3.1, §7.1 P4, D10).
The :mod:`agent_team.transport.claude_code_adapter` module ships the §3.3.1
transport contract with a dependency-injected ``delivery`` seam: the adapter
renders the question-set into a prompt body and hands it to a sink called as
``delivery(session_hint=..., prompt=...) -> str`` whose job is to surface the
prompt inside a Claude-Code session and return the Claude **session id** used as
the ``channel_ref`` locator. The foundation's default sink refuses to act so
nothing ships provisioned; this module supplies the **production** sink, backed
by a local **file drop** the Mac harness reads, that the P4 (Claude-Code) live
wiring injects.
Channel model (file drop)
-------------------------
The R720 box has no standing write path into Adam's interactive Claude-Code
session, so delivery is **SSH-invoked from the Mac side** (D10). The live sink
writes the rendered prompt to a drop directory on the Mac filesystem; the
Claude-Code harness polls that directory, surfaces the prompt inline, and Adam
answers it. The sink returns the Claude **session id** (the file stem, derived
from the embedded ``question_id``) which §3.3.1 names as this transport's
``channel_ref`` locator. The answer travels back the same way: the harness drops
an answer file the box reads and feeds to
:meth:`ClaudeCodeAdapter.parse_answer`.
Deferred / injected I/O (mirrors :func:`agent_team.graph.build_sqlite_checkpointer`)
-----------------------------------------------------------------------------------
All filesystem access is injected or deferred so this module imports cleanly in
pre-deploy / test environments and is unit-testable with fakes:
* the prompt-writer and answer-reader are injected callables (``writer`` /
``reader``) defaulting to thin wrappers over the local filesystem; and
* the standard-library ``pathlib`` import is the only hard dependency — there is
no SDK, network, SSH, or secret touched here. A missing / unwritable drop
directory raises a clear :class:`ClaudeCodeDeliveryError` so a misconfigured
deploy fails loudly rather than silently reporting a delivery that never
reached Adam.
Scope (P2 default; P3 inert)
----------------------------
This is the human-in-the-loop *question delivery* path used by the production
P2 (clarify -> plan -> review) pipeline. It performs no CI, OIDC, git/patch
apply, or network calls; the P3 builder/verifier apply-verify workflow is held
for a separate review gate and is neither shipped nor enabled here.
"""
from __future__ import annotations
import os
from pathlib import Path
from typing import Any, Callable
from agent_team.transport.claude_code_adapter import (
ClaudeCodeAdapter,
ClaudeCodeDeliveryError,
parse_channel_ref,
)
__all__ = [
"PROMPT_SUFFIX",
"build_claude_code_delivery",
"build_file_drop_reader",
"build_file_drop_writer",
"build_live_claude_code_transport",
"read_answer_payload",
]
# File extension for a dropped prompt file. The Mac harness globs this suffix in
# the drop directory to discover pending questions.
PROMPT_SUFFIX = ".prompt.txt"
# File extension for a dropped answer file written back by the Mac harness. The
# box globs this suffix to discover answered questions.
ANSWER_SUFFIX = ".answer.txt"
# Type of the injected prompt-writer: surfaces ``prompt`` at ``path``.
PromptWriter = Callable[[Path, str], None]
# Type of the injected answer-reader: returns the answer text at ``path`` (or
# ``None`` if the harness has not dropped an answer yet).
AnswerReader = Callable[[Path], "str | None"]
def _drop_path(drop_dir: Path, session_id: str, suffix: str) -> Path:
"""Resolve the on-disk path for a drop file.
The ``session_id`` is the file stem, so a prompt and its answer share a stem
and differ only by suffix. ``session_id`` is sanitized to a single path
component (no separators) so a hostile/garbled id cannot escape ``drop_dir``.
"""
# Replace os.sep + os.altsep AND an explicit backslash: on POSIX os.altsep
# is None, so a literal backslash would otherwise survive (defense-in-depth,
# even though stripping "/" already prevents traversal on the deploy targets).
safe = session_id.strip().replace(os.sep, "_").replace("\\", "_")
if os.altsep:
safe = safe.replace(os.altsep, "_")
safe = safe.lstrip(".") or "session"
return drop_dir / f"{safe}{suffix}"
def build_file_drop_writer(drop_dir: Path | str) -> PromptWriter:
"""Build the default filesystem prompt-writer for the live sink.
The returned callable writes ``prompt`` to ``path`` (UTF-8), creating the
drop directory if needed. Filesystem access is deferred to call time so this
factory is side-effect-free at import. A write failure (unwritable / missing
parent) propagates so the caller can treat the post as failed.
"""
base = Path(drop_dir)
def _writer(path: Path, prompt: str) -> None:
base.mkdir(parents=True, exist_ok=True)
path.write_text(prompt, encoding="utf-8")
return _writer
def build_file_drop_reader(drop_dir: Path | str) -> AnswerReader:
"""Build the default filesystem answer-reader for the inbound path.
The returned callable reads the answer text at ``path`` (UTF-8), returning
``None`` when the harness has not yet dropped an answer file. Filesystem
access is deferred to call time.
"""
Path(drop_dir) # validate / normalize eagerly; read happens at call time.
def _reader(path: Path) -> str | None:
try:
return path.read_text(encoding="utf-8")
except FileNotFoundError:
return None
return _reader
def build_claude_code_delivery(
drop_dir: Path | str,
*,
writer: PromptWriter | None = None,
) -> Callable[..., str]:
"""Build a live file-drop ``delivery`` sink (§3.3.1, P4, D10).
The returned callable matches the adapter's injected-sink contract,
``delivery(session_hint=..., prompt=...) -> str``: it writes the rendered
``prompt`` to a file in ``drop_dir`` for the Mac harness to read and returns
the Claude **session id** the adapter folds into the ``channel_ref``.
The session id is derived from the ``question_id`` embedded in ``prompt`` via
the shared ``<!-- shq:<question_id> -->`` marker (so the prompt and its
answer share a deterministic file stem), falling back to ``session_hint``
when present. A prompt with neither an embedded marker nor a ``session_hint``
cannot be addressed, so the sink raises :class:`ClaudeCodeDeliveryError`
rather than dropping an unaddressable file.
``writer`` (optional) injects the prompt-writer for testability; any
``Callable[[Path, str], None]`` works. When omitted, a filesystem writer over
``drop_dir`` is built via :func:`build_file_drop_writer`. A write failure is
wrapped in :class:`ClaudeCodeDeliveryError` so a failed drop leaves the
ledger row ``open`` (no ``channel_ref``) for idempotent reconcile retry,
mirroring a failed Slack post.
"""
base = Path(drop_dir)
write = writer if writer is not None else build_file_drop_writer(base)
def _delivery(*, session_hint: str, prompt: str) -> str:
session_id = _session_id_for(prompt=prompt, session_hint=session_hint)
path = _drop_path(base, session_id, PROMPT_SUFFIX)
try:
write(path, prompt)
except OSError as exc:
raise ClaudeCodeDeliveryError(
f"failed to write Claude-Code prompt drop at {path!s}: {exc}"
) from exc
return session_id
return _delivery
def build_live_claude_code_transport(
drop_dir: Path | str,
*,
session_hint: str = "",
writer: PromptWriter | None = None,
) -> ClaudeCodeAdapter:
"""Build a :class:`ClaudeCodeAdapter` wired to a live file-drop sink.
Convenience constructor for the P4 live coordinator: equivalent to
``ClaudeCodeAdapter(build_claude_code_delivery(drop_dir, writer=writer),
session_hint=session_hint)``. See :func:`build_claude_code_delivery` for the
drop-directory / writer / session-id semantics.
"""
return ClaudeCodeAdapter(
build_claude_code_delivery(drop_dir, writer=writer),
session_hint=session_hint,
)
def read_answer_payload(
channel_ref: str,
drop_dir: Path | str,
*,
reader: AnswerReader | None = None,
) -> dict[str, Any] | None:
"""Read a dropped answer for ``channel_ref`` into a ``parse_answer`` payload.
Resolves the answer file (the prompt's stem + :data:`ANSWER_SUFFIX`) for the
Claude session in ``channel_ref``, reads it via the injected ``reader``
(defaulting to a filesystem reader over ``drop_dir``), and returns a payload
dict ready for :meth:`ClaudeCodeAdapter.parse_answer`. The payload echoes the
original ``channel_ref`` so the ``question_id`` round-trips from the ref
alone even if the answer text carries no marker.
Returns ``None`` when no answer has been dropped yet (the harness has not
answered), so the reconcile loop can poll idempotently. Raises
:class:`ValueError` for a ``channel_ref`` that is not a Claude-Code ref, so a
caller cannot silently read the wrong transport's drop.
"""
parsed = parse_channel_ref(channel_ref)
if parsed is None:
raise ValueError(
f"not a Claude-Code channel_ref; cannot read answer drop: {channel_ref!r}"
)
session_id, _question_id = parsed
base = Path(drop_dir)
read = reader if reader is not None else build_file_drop_reader(base)
path = _drop_path(base, session_id, ANSWER_SUFFIX)
answer_text = read(path)
if answer_text is None:
return None
return {"channel_ref": channel_ref, "answer": answer_text}
def _session_id_for(*, prompt: str, session_hint: str) -> str:
"""Derive the file-stem session id for a prompt drop.
Prefers the ``question_id`` embedded in ``prompt`` via the shared marker (so
the prompt and its answer share a deterministic stem), then a non-blank
``session_hint``. Raises :class:`ClaudeCodeDeliveryError` when neither is
available, since an unaddressable drop would be unrecoverable.
"""
embedded = _embedded_question_id(prompt)
if embedded:
return embedded
hint = session_hint.strip()
if hint:
return hint
raise ClaudeCodeDeliveryError(
"cannot address a Claude-Code prompt drop: no embedded question_id marker "
"and no session_hint"
)
def _embedded_question_id(prompt: str) -> str | None:
"""Recover the ``question_id`` embedded in a rendered prompt body, or ``None``.
Delegates to the adapter's own marker parsing through a throwaway payload so
the marker grammar stays single-sourced in
:mod:`agent_team.transport.claude_code_adapter`.
"""
adapter = ClaudeCodeAdapter()
try:
question_id, _answer, _via = adapter.parse_answer({"prompt": prompt})
except ValueError:
return None
return question_id

View file

@ -0,0 +1,293 @@
"""GitHub-issue INTAKE poller: a labeled issue becomes a pipeline task (§3.3.1).
This is the *inbound front door* for the GitHub transport. Where
:mod:`agent_team.transport.github_adapter` delivers clarifier question-sets
*outbound* (and parses answers back), this leaf runs the other direction: it
polls a repository for open issues carrying a configured label and turns each
not-yet-ingested issue into one pipeline task by calling the coordinator's
intake entry,
:meth:`agent_team.coordinator.Coordinator.start_task` (``task_text=<issue
title+body>``, ``transport_name="github"``).
The shape mirrors the slack_listener seam: everything network/SDK is
**injected** so the poller is fully unit-testable with no GitHub SDK and no
socket:
* ``client`` is a small :class:`GithubIssueClient` protocol:
``list_open_issues(label) -> iterable of issue mappings``. Production wires a
thin client over the GitHub REST API (deferred import, see
:func:`build_default_issue_client`); tests pass an in-memory fake.
* ``coordinator`` is anything exposing ``start_task(task_text=...,
transport_name=...)``: the live :class:`~agent_team.coordinator.Coordinator`
in production, a stub in tests. No model or transport is touched here.
De-duplication (P3 scope note):
The poller tracks already-ingested issue ids in an **in-memory** set, so a
re-poll over the same open issue does not start a second task. This is
deliberately simple for now: it does NOT survive a process restart. Durable
de-dup (a ledger table of ingested issue ids, mirroring the
``pending_questions`` discipline) is a FOLLOW-UP and is intentionally not
shipped here. After a restart an already-ingested-but-still-open issue would
be re-ingested; document that and treat the in-memory set as a best-effort
guard, not a durable contract.
Design constraints honoured here (pre-deployment scaffolding):
* **No live infrastructure.** Nothing is provisioned or called at import.
The GitHub client is dependency-injected; the default client's SDK/HTTP
import is DEFERRED (mirrors
:func:`agent_team.graph.build_sqlite_checkpointer` and the slack_listener
SDK discipline), so this module imports cleanly with no optional SDK
present and the unit tests stay fully hermetic.
* **P2 stays the production default; this is OPT-IN and INERT.** This module
does no CI, no OIDC, no git/patch apply, and no network to GitHub Actions.
It only reads issues and calls the coordinator's existing intake entry.
* **Secrets never committed.** The default client reads the GitHub token
from the environment at call time, never from source.
"""
from __future__ import annotations
import logging
from typing import Any, Iterable, Protocol
__all__ = [
"GithubIntake",
"GithubIssueClient",
"build_default_issue_client",
"issue_task_text",
]
_LOG = logging.getLogger(__name__)
# The transport name handed to the coordinator's intake entry so the resulting
# task's clarifier question-sets route over the GitHub adapter (§3.3.1 D10).
GITHUB_TRANSPORT_NAME = "github"
# Environment variable the default client reads the GitHub token from at call
# time (never stored in source/state). Mirrors the github_adapter default.
DEFAULT_TOKEN_ENV = "GITHUB_TOKEN"
class GithubIssueClient(Protocol):
"""Injected GitHub issue source: list open issues carrying a label.
A narrow read-only seam so the poller has no hard dependency on any GitHub
SDK and the tests pass a pure in-memory fake. Each returned issue is a
mapping with at least an ``id`` (or ``number``) and a ``title``; ``body`` is
optional. The mapping shape mirrors the GitHub REST issue object so the
production client can return the API JSON unchanged.
"""
def list_open_issues(self, *, label: str) -> Iterable[dict[str, Any]]:
"""Return the open issues carrying ``label`` (most-recent-first is fine)."""
...
def issue_task_text(issue: dict[str, Any]) -> str:
"""Render one issue's intake ``task_text`` from its title + body.
The pipeline's task description is the issue title followed by its body (a
blank line between them when both are present). A missing/empty body yields
just the title; a missing/empty title falls back to ``issue #<id>`` so the
task is never an empty string. Whitespace is stripped at the edges so a
trailing-newline body does not produce trailing blank lines.
"""
title = str(issue.get("title") or "").strip()
body = str(issue.get("body") or "").strip()
if not title:
title = f"issue #{_issue_id(issue)}"
if body:
return f"{title}\n\n{body}"
return title
def _issue_id(issue: dict[str, Any]) -> str:
"""Return the de-dup identity for ``issue`` as a string.
Prefers the GitHub global ``id`` (stable across renames); falls back to the
per-repo ``number`` when ``id`` is absent (some payload shapes / fakes carry
only ``number``). Stringified so heterogeneous int/str ids compare cleanly
in the ingested set.
"""
raw = issue.get("id")
if raw is None:
raw = issue.get("number")
return str(raw)
class GithubIntake:
"""Poll a repo for labeled issues and start one pipeline task per new issue.
Construct with an injected ``client`` (a :class:`GithubIssueClient`), an
injected ``coordinator`` (anything exposing
``start_task(task_text=..., transport_name=...)``), and the ``label`` that
flags an issue as pipeline intake. Call :meth:`poll_once` on a cadence (an
operator loop or a cron); each call lists the open labeled issues and starts
a task for every one not yet ingested.
De-dup is in-memory only (see the module docstring): the set of ingested
issue ids lives on the instance, so a re-poll within one process never
double-ingests, but a restart loses the set. Durable de-dup is a follow-up.
Nothing here touches the network or any SDK directly (the client does, and
it is injected), so the whole poller is unit-testable with a fake client and
a stub coordinator.
"""
def __init__(
self,
*,
client: GithubIssueClient,
coordinator: Any,
label: str,
) -> None:
"""Bind the poller to one client, coordinator, and intake label.
Args:
client: The injected issue source. Its ``list_open_issues`` is the
only GitHub call the poller makes.
coordinator: The intake target. Must expose
``start_task(task_text=..., transport_name=...)``: the live
:class:`~agent_team.coordinator.Coordinator` in production.
label: The issue label that marks an issue as pipeline intake. Only
issues the client returns for this label are considered; an
empty label is rejected so a misconfiguration cannot ingest
every open issue.
"""
if not label:
raise ValueError(
"GithubIntake requires a non-empty intake label; an empty label "
"would ingest every open issue"
)
self._client = client
self._coordinator = coordinator
self._label = label
# In-memory de-dup set (P3 scope: best-effort, NOT durable across a
# restart; see the module docstring). Tracks issue ids already turned
# into tasks so a re-poll does not double-ingest.
self._ingested: set[str] = set()
@property
def label(self) -> str:
"""The configured intake label (read-only)."""
return self._label
@property
def ingested_ids(self) -> frozenset[str]:
"""A snapshot of the issue ids already ingested this process (read-only)."""
return frozenset(self._ingested)
def poll_once(self) -> list[str]:
"""List the labeled open issues and start a task for each new one.
One maintenance pass:
1. Ask the injected client for the open issues carrying the configured
label (:meth:`GithubIssueClient.list_open_issues`).
2. For each issue NOT already in the in-memory ingested set, call
``coordinator.start_task(task_text=<title+body>,
transport_name="github")`` and record its id so a subsequent poll
does not re-ingest it.
Issues already ingested this process are skipped (the in-memory de-dup),
and any issue the client returns without the label is *not* expected
(the client filters by label) but is ignored defensively if present.
An issue id is recorded as ingested ONLY after ``start_task`` returns,
so a failing intake leaves the issue eligible for retry on the next poll
rather than silently dropping it.
Returns the list of issue ids ingested on THIS pass (empty when nothing
new), so an operator loop can log/meter intake volume.
"""
ingested_now: list[str] = []
for issue in self._client.list_open_issues(label=self._label):
issue_id = _issue_id(issue)
if issue_id in self._ingested:
_LOG.debug("github-intake: issue %s already ingested; skip", issue_id)
continue
task_text = issue_task_text(issue)
_LOG.info(
"github-intake: starting task for issue %s (label=%s)",
issue_id,
self._label,
)
# start_task is the committed coordinator intake entry; the resulting
# task's clarifier question-sets route over the GitHub adapter. Record
# the id only after the call returns so a raise leaves the issue
# eligible for retry on the next poll (no silent drop).
self._coordinator.start_task(
task_text=task_text,
transport_name=GITHUB_TRANSPORT_NAME,
)
self._ingested.add(issue_id)
ingested_now.append(issue_id)
return ingested_now
def build_default_issue_client(
*,
owner: str,
repo: str,
token_env: str = DEFAULT_TOKEN_ENV,
api_root: str = "https://api.github.com",
) -> GithubIssueClient:
"""Build the production read-only issue client (deferred SDK/HTTP import).
Returns a :class:`GithubIssueClient` that lists a repo's open issues by
label over the GitHub REST API. The HTTP machinery (``urllib``) and the
token read are deferred to call time (mirroring
:func:`agent_team.graph.build_sqlite_checkpointer` and the slack_listener
SDK discipline), so importing this module never touches the network and the
unit tests (which inject a fake client) never reach this path.
The token is read from ``token_env`` at call time and sent as a bearer
credential; it is never stored in source or logged. ``api_root`` is
overridable for GitHub Enterprise.
This is intentionally a thin, read-only lister: it issues a single GET to
the issues endpoint with ``state=open&labels=<label>`` and returns the
parsed JSON array unchanged (each element is a GitHub issue object, which
already carries ``id`` / ``number`` / ``title`` / ``body``). It performs no
CI, OIDC, write, or GitHub-Actions call; it only reads issues.
"""
class _RestIssueClient:
"""Stdlib-only GitHub REST issue lister (built lazily, no import-time HTTP)."""
def __init__(self) -> None:
self._owner = owner
self._repo = repo
self._token_env = token_env
self._api_root = api_root.rstrip("/")
def list_open_issues(self, *, label: str) -> Iterable[dict[str, Any]]:
import json
import os
from urllib import parse as _urlparse
from urllib import request as _urlrequest
token = os.environ.get(self._token_env)
if not token:
raise RuntimeError(
f"no GitHub token available (env {self._token_env!r} unset); "
"cannot list issues for intake"
)
query = _urlparse.urlencode({"state": "open", "labels": label})
url = f"{self._api_root}/repos/{self._owner}/{self._repo}/issues?{query}"
request = _urlrequest.Request(url, method="GET")
request.add_header("Authorization", f"Bearer {token}")
request.add_header("Accept", "application/vnd.github+json")
request.add_header("X-GitHub-Api-Version", "2022-11-28")
with _urlrequest.urlopen(request) as response: # noqa: S310 (trusted api host)
raw = response.read().decode("utf-8")
data = json.loads(raw) if raw else []
# The issues endpoint can include pull requests (they share the
# endpoint); filter them out so a PR is never intaken as an issue.
return [item for item in data if "pull_request" not in item]
return _RestIssueClient()

View file

@ -0,0 +1,229 @@
"""Live ``requests``-backed GitHub poster (design §3.3.1, §7.1 P4 — GitHub).
The :mod:`agent_team.transport.github_adapter` module ships the §3.3.1 transport
contract with a dependency-injected ``http_post`` seam: the adapter renders the
question-set into a Markdown comment body (embedding the
``<!-- shq:<question_id> -->`` marker so an inbound answer maps back) and hands
the REST POST to an
``HttpPost = (url, *, headers, json_body) -> (status, data)`` whose job is to
perform the real ``POST /repos/{owner}/{repo}/issues/{n}/comments`` and return
the GitHub comment payload carrying the new comment ``id``. The foundation's
default ``http_post`` is a stdlib-only (``urllib``) poster invoked only on an
actual delivery, so nothing ships provisioned; this module supplies the
**production** poster, backed by a thin ``requests`` session, that the P4 (GitHub)
live wiring injects.
Why a thin ``requests`` shim (and not PyGithub):
The adapter already speaks the GitHub REST API directly — it builds the
comments URL, the ``Authorization: Bearer`` / ``X-GitHub-Api-Version``
headers, and the ``{"body": ...}`` JSON itself, then hands a plain
``(url, headers, json_body)`` POST to the seam. The live poster therefore
only needs a minimal HTTP client, not the full PyGithub object model. A
``requests`` session keeps the shim small and fully mirrors the seam the
adapter already accepts.
Deferred import (mirrors :func:`agent_team.graph.build_sqlite_checkpointer` and
:func:`agent_team.transport.slack_live.build_slack_poster`):
``requests`` is an optional dependency that may be absent in pre-deploy / test
environments, so this module imports cleanly without it. The import is deferred
to the moment a live client is actually constructed, and a missing package
raises a clear :class:`RuntimeError` so a misconfigured deploy fails loudly
rather than silently. The token is likewise resolved at build time (falling back
to ``GITHUB_TOKEN``); a missing token raises a clear :class:`RuntimeError`.
The marker round-trip is owned by the adapter, not the poster: the poster is the
pure network seam. :meth:`GitHubTransport.post_question` embeds the
``question_id`` marker in the comment body it hands to this poster, and
:meth:`GitHubTransport.parse_answer` recovers it from an inbound reply, so the
``question_id`` survives end-to-end without the poster needing to know about it.
Production-default safety (P2 stays clarify->plan->review): this module is live
**transport I/O only** for the human gate. It performs exactly one operation —
post an issue/PR comment — and contains NO CI calls, NO OIDC, NO git/patch
apply, and NO GitHub Actions network. The P3 builders/verifier live apply/verify
workflow is held for a separate review gate and is deliberately absent here.
"""
from __future__ import annotations
import os
from typing import Any
from agent_team.transport.github_adapter import (
GITHUB_API_ROOT,
GitHubApiError,
GitHubTransport,
HttpPost,
)
__all__ = [
"build_github_poster",
"build_live_github_transport",
]
def build_github_poster(token: str | None = None, *, client: Any = None) -> HttpPost:
"""Build a live ``requests``-backed :data:`HttpPost` seam (§3.3.1, P4).
The returned callable implements the adapter's
``(url, *, headers, json_body) -> (status, data)`` seam: it performs the
real ``POST`` against the GitHub REST API and returns the response status
plus parsed JSON body so :meth:`GitHubTransport.post_question` can read the
new comment ``id`` as the ``channel_ref`` the ledger stores. A non-2xx
response is surfaced as :class:`GitHubApiError` (the adapter also guards the
status, but the poster raises early so a failed POST never looks like a
success with an empty body).
``client`` (optional) injects a pre-built HTTP client for testability; any
object exposing ``post(url, *, headers, json)`` and returning a response
with ``status_code`` plus a ``json()`` method works (the ``requests``
``Session`` shape). When omitted, a ``requests.Session`` is constructed
lazily from ``token`` (falling back to the ``GITHUB_TOKEN`` environment
variable). The ``requests`` import is deferred so this module imports
cleanly without the optional package; a missing package or a missing token
raises a clear :class:`RuntimeError`.
The poster is the pure network seam: the ``question_id`` marker is embedded
by the adapter in the comment body it passes through ``json_body``, so the
posted comment carries the marker without the poster handling it. See the
module docstring for the full rationale.
"""
if client is None:
client = _build_session(token)
def _poster(
url: str,
*,
headers: dict[str, str],
json_body: dict[str, Any],
) -> tuple[int, dict[str, Any]]:
response = client.post(url, headers=headers, json=json_body)
status = _status_of(response)
data = _json_of(response)
if not (200 <= status < 300):
raise GitHubApiError(status, _body_repr(data))
return status, data
return _poster
def build_live_github_transport(
*,
owner: str,
repo: str,
issue_number: int,
token: str | None = None,
client: Any = None,
api_root: str = GITHUB_API_ROOT,
) -> GitHubTransport:
"""Build a :class:`GitHubTransport` wired to a live ``requests`` poster.
Convenience constructor for the P4 live coordinator: builds the live poster
and binds it to the issue (or PR) thread where the human gate posts its
question-set comment and reads the reply. See :func:`build_github_poster`
for the token / client / deferred-import semantics.
The transport builds the ``Authorization`` header itself (from a
``token_provider``), so the resolved token is threaded in once and shared
with the poster's session: a single token source backs the whole live path.
When neither ``token`` nor ``client`` is given, the token still resolves from
``GITHUB_TOKEN`` at call time so the deploy fails loudly if it is unset.
``api_root`` is forwarded to the transport so a GitHub Enterprise host can be
targeted; the poster itself is endpoint-agnostic (the adapter builds the URL).
"""
resolved = _resolve_token(token)
return GitHubTransport(
owner=owner,
repo=repo,
issue_number=issue_number,
http_post=build_github_poster(resolved, client=client),
api_root=api_root,
token_provider=(lambda: resolved) if resolved is not None else None,
)
def _build_session(token: str | None) -> Any:
"""Lazily construct a ``requests.Session`` (deferred optional import).
Raises a clear :class:`RuntimeError` if ``requests`` is not installed or no
token is resolvable (neither ``token`` nor ``GITHUB_TOKEN``), so a
misconfigured deploy fails loudly rather than silently. The token is held on
the session purely so the same authenticated client can be reused; the
adapter still sets the ``Authorization`` header per call.
"""
try:
import requests
except ImportError as exc: # pragma: no cover - depends on optional dep
raise RuntimeError(
"requests is unavailable; install the 'requests' package to build "
"a live GitHub poster (P4), or inject a 'client' for testing."
) from exc
resolved = _resolve_token(token)
if not resolved:
raise RuntimeError(
"No GitHub token available; pass 'token' or set the GITHUB_TOKEN "
"environment variable to build a live GitHub poster."
)
session = requests.Session()
session.headers.update({"Authorization": f"Bearer {resolved}"})
return session
def _resolve_token(token: str | None) -> str | None:
"""Resolve the GitHub token, falling back to ``GITHUB_TOKEN`` (read at call time).
Returns ``None`` when no token is available so callers can decide whether a
missing token is fatal (the poster path) or deferrable (an injected-client
test path that never needs one).
"""
return token or os.environ.get("GITHUB_TOKEN")
def _status_of(response: Any) -> int:
"""Read the HTTP status code from a ``requests``-like response.
``requests.Response`` exposes ``status_code``; a test double may instead
expose ``status``. Either is accepted so the poster stays usable with a
minimal fake.
"""
for attr in ("status_code", "status"):
value = getattr(response, attr, None)
if value is not None:
return int(value)
raise TypeError(
"GitHub response exposes no 'status_code'/'status'; expected a "
f"requests-like response, got {type(response)!r}"
)
def _json_of(response: Any) -> dict[str, Any]:
"""Parse the JSON body of a ``requests``-like response to a dict.
A successful create returns the new comment object (carrying ``id``); an
empty body coerces to ``{}`` so the adapter's missing-id guard fires with a
clear message rather than an attribute error.
"""
parser = getattr(response, "json", None)
if not callable(parser):
raise TypeError(
"GitHub response exposes no callable 'json()'; expected a "
f"requests-like response, got {type(response)!r}"
)
data = parser()
return data if isinstance(data, dict) else {}
def _body_repr(data: dict[str, Any]) -> str:
"""Render a response body for a :class:`GitHubApiError` message.
Best-effort JSON; falls back to ``repr`` so the error is always constructible
even for an exotic body.
"""
try:
import json
return json.dumps(data)
except (TypeError, ValueError): # pragma: no cover - exotic body
return repr(data)

View file

@ -38,6 +38,9 @@ Subcommands (P1 surface):
``--confirm``; audit-logged). Maps to :func:`answer_question`. ``--confirm``; audit-logged). Maps to :func:`answer_question`.
* ``supersede`` — mark a stale question ``superseded`` (DESTRUCTIVE: requires * ``supersede`` — mark a stale question ``superseded`` (DESTRUCTIVE: requires
``--confirm``; audit-logged). Maps to :func:`supersede_question`. ``--confirm``; audit-logged). Maps to :func:`supersede_question`.
* ``intake-github`` — poll a repo for labeled open issues and start one pipeline
task per not-yet-ingested issue (one pass). Reads issues + calls the committed
coordinator intake entry only; no CI, OIDC, or git/patch apply. Opt-in/inert.
Exit codes: ``0`` success, ``1`` operational failure (e.g. row not found, the Exit codes: ``0`` success, ``1`` operational failure (e.g. row not found, the
compare-and-set lost the race), ``2`` usage error (argparse). compare-and-set lost the race), ``2`` usage error (argparse).
@ -74,9 +77,12 @@ from agent_team.db.schema import ( # noqa: E402 (path bootstrap must precede)
) )
from agent_team.transport.base import Transport # noqa: E402 (path bootstrap) from agent_team.transport.base import Transport # noqa: E402 (path bootstrap)
# Transport choices the start/serve commands accept (§3.3.1 D10). Only ``slack`` # Transport choices the start/serve/intake commands accept (§3.3.1 D10). All
# has a live adapter wired for the P1 CLI; the others are accepted for forward # three now have a live human-gate adapter wired in _build_transport: ``slack``
# compatibility and gated in _build_transport. # (slack_sdk), ``github`` (requests issue/PR comment), and ``claude_code`` (local
# file drop). Each builder defers its optional SDK / token resolution to call
# time, so a missing dep/credential fails loudly only when that transport is
# actually selected.
_TRANSPORT_CHOICES: tuple[str, ...] = ("slack", "github", "claude_code") _TRANSPORT_CHOICES: tuple[str, ...] = ("slack", "github", "claude_code")
__all__ = [ __all__ = [
@ -511,10 +517,23 @@ def _build_transport(args: argparse.Namespace) -> Any:
"""Build the transport for a coordinator command (lazy; token-tolerant). """Build the transport for a coordinator command (lazy; token-tolerant).
``--dry-run`` (or any transport in dry-run) yields a non-posting transport so ``--dry-run`` (or any transport in dry-run) yields a non-posting transport so
intake works without credentials. Otherwise the live Slack transport is intake works without credentials. Otherwise the live transport is built
constructed lazily from ``SLACK_BOT_TOKEN`` / ``SLACK_CHANNEL``; GitHub and lazily from the environment for the chosen ``--transport`` (so import,
Claude-Code live transports are not wired for the P1 CLI surface and raise a ``--help``, and ledger commands never need a token):
clear error rather than pretending to post.
* ``slack`` -> :func:`build_live_slack_transport` over ``SLACK_BOT_TOKEN`` /
``SLACK_CHANNEL``.
* ``github`` -> :func:`build_live_github_transport` over ``GITHUB_TOKEN`` and
the issue thread ``GITHUB_OWNER`` / ``GITHUB_REPO`` /
``GITHUB_ISSUE_NUMBER``. This is the §3.3.1 human-gate I/O only (post an
issue/PR comment); it carries NO CI, OIDC, git/patch apply, or GitHub
Actions network (the P3 apply/verify workflow is held for the security
gate).
* ``claude_code`` -> :func:`build_live_claude_code_transport` over the local
file-drop directory ``CLAUDE_CODE_DROP_DIR`` the Mac harness polls (D10).
Each live builder defers its optional SDK / token resolution to call time, so
a missing dependency or credential fails loudly here rather than at import.
""" """
if getattr(args, "dry_run", False): if getattr(args, "dry_run", False):
return _DryRunTransport() return _DryRunTransport()
@ -523,12 +542,74 @@ def _build_transport(args: argparse.Namespace) -> Any:
channel = os.environ.get("SLACK_CHANNEL", "") channel = os.environ.get("SLACK_CHANNEL", "")
return build_live_slack_transport(channel) return build_live_slack_transport(channel)
if args.transport == "github":
from agent_team.transport.github_live import build_live_github_transport
owner, repo, issue_number = _github_thread_from_env()
# Token resolves from GITHUB_TOKEN inside the builder (call-time read);
# a missing token fails loudly there rather than being captured here.
return build_live_github_transport(
owner=owner,
repo=repo,
issue_number=issue_number,
token=os.environ.get("GITHUB_TOKEN") or None,
)
if args.transport == "claude_code":
from agent_team.transport.claude_code_live import (
build_live_claude_code_transport,
)
drop_dir = os.environ.get("CLAUDE_CODE_DROP_DIR", "")
if not drop_dir:
raise SystemExit(
"live transport 'claude_code' requires CLAUDE_CODE_DROP_DIR (the "
"Mac file-drop directory the Claude-Code harness polls); set it, "
"or use --dry-run for a no-token dry run"
)
return build_live_claude_code_transport(drop_dir)
raise SystemExit( raise SystemExit(
f"live transport '{args.transport}' is not wired for the run-team CLI; " f"live transport '{args.transport}' is not wired for the run-team CLI; "
"use --transport slack, or --dry-run for a no-token dry run" "use --transport slack/github/claude_code, or --dry-run for a no-token "
"dry run"
) )
def _github_thread_from_env() -> tuple[str, str, int]:
"""Resolve the GitHub issue thread (owner/repo/issue) from the environment.
The live GitHub transport posts the clarifier question-set as a comment on a
fixed ``owner/repo#issue_number`` thread, so the thread is configured via
``GITHUB_OWNER`` / ``GITHUB_REPO`` / ``GITHUB_ISSUE_NUMBER`` (read lazily so
the value is never captured at import). A missing or non-integer value raises
a clear :class:`SystemExit` rather than building a half-configured transport.
"""
owner = os.environ.get("GITHUB_OWNER", "")
repo = os.environ.get("GITHUB_REPO", "")
raw_issue = os.environ.get("GITHUB_ISSUE_NUMBER", "")
missing = [
name
for name, value in (
("GITHUB_OWNER", owner),
("GITHUB_REPO", repo),
("GITHUB_ISSUE_NUMBER", raw_issue),
)
if not value
]
if missing:
raise SystemExit(
"live transport 'github' requires "
f"{', '.join(missing)}; set the issue thread (owner/repo/issue), or "
"use --dry-run for a no-token dry run"
)
try:
issue_number = int(raw_issue)
except ValueError:
raise SystemExit(
f"GITHUB_ISSUE_NUMBER must be an integer, got {raw_issue!r}"
) from None
return owner, repo, issue_number
class _DryRunTransport(Transport): class _DryRunTransport(Transport):
"""A non-posting transport for ``--dry-run`` intake (no token, no Slack). """A non-posting transport for ``--dry-run`` intake (no token, no Slack).
@ -591,6 +672,46 @@ def _cmd_serve(args: argparse.Namespace, *, out: Any) -> int:
return 0 # pragma: no cover - serve() loops until interrupted return 0 # pragma: no cover - serve() loops until interrupted
def _cmd_intake_github(args: argparse.Namespace, *, out: Any) -> int:
"""Poll a repo for labeled issues and start one task per new issue (§3.3.1).
The GitHub-issue intake front door: builds a :class:`Coordinator` (transport
from the lazy factory; ``--dry-run`` posts nowhere), runs ``setup``, then
constructs a :class:`agent_team.transport.github_intake.GithubIntake` over a
read-only REST issue client
(:func:`~agent_team.transport.github_intake.build_default_issue_client`,
deferred-import, reads ``GITHUB_TOKEN`` at call time) and runs ONE
:meth:`~agent_team.transport.github_intake.GithubIntake.poll_once`. An
operator (or a cron) re-runs the command on a cadence; de-dup is in-memory
per process, so each run is a single pass.
OPT-IN and INERT: this only reads labeled issues and calls the committed
coordinator intake entry — no CI, OIDC, git/patch apply, or GitHub-Actions
network. Owner/repo/label come from CLI flags; the token comes from
``GITHUB_TOKEN``. Prints the issue ids ingested on this pass (one per line).
"""
from agent_team.transport.github_intake import (
GithubIntake,
build_default_issue_client,
)
coordinator = _build_coordinator(args)
coordinator.setup()
client = build_default_issue_client(owner=args.owner, repo=args.repo)
intake = GithubIntake(
client=client,
coordinator=coordinator,
label=args.label,
)
ingested = intake.poll_once()
for issue_id in ingested:
print(issue_id, file=out)
if not ingested:
print("github-intake: no new labeled issues to ingest", file=sys.stderr)
return 0
def _cmd_force_resume(args: argparse.Namespace, *, out: Any) -> int: def _cmd_force_resume(args: argparse.Namespace, *, out: Any) -> int:
"""Force-resume a parked task's question (destructive; audit-logged). """Force-resume a parked task's question (destructive; audit-logged).
@ -828,6 +949,39 @@ def build_parser() -> argparse.ArgumentParser:
) )
p_serve.set_defaults(func=_cmd_serve) p_serve.set_defaults(func=_cmd_serve)
p_intake = sub.add_parser(
"intake-github",
help="poll a repo for labeled issues and start one task per new issue",
)
p_intake.add_argument(
"--owner",
required=True,
help="GitHub repository owner / org login to poll for intake issues",
)
p_intake.add_argument(
"--repo",
required=True,
help="GitHub repository name to poll for intake issues",
)
p_intake.add_argument(
"--label",
required=True,
help="issue label that flags an issue as pipeline intake (non-empty)",
)
p_intake.add_argument(
"--transport",
choices=_TRANSPORT_CHOICES,
default="github",
help="channel for delivering clarifier questions (default: github)",
)
p_intake.add_argument(
"--dry-run",
action="store_true",
dest="dry_run",
help="use a non-posting transport (no token needed; ingest still runs)",
)
p_intake.set_defaults(func=_cmd_intake_github)
return parser return parser

View file

@ -0,0 +1,375 @@
"""Unit tests for agent_team.nodes.build_verify_subgraph (P3-INERT topology).
These tests prove the build -> verify subgraph TOPOLOGY is correctly inert:
* the BUILD node proposes a candidate diff via an INJECTED fake builder and
advances to VERIFY;
* the VERIFY node, fed a fake authenticated-pass ``ci_result``, routes to the
approved / PR terminus;
* the VERIFY node with the DEFAULT (None) fetcher — and with a failing fetcher —
BLOCKs and routes to PARKED, never fabricating a pass;
* an LLM fix-proposal can NEVER flip a failing verdict to pass (the gate is the
sole pass authority).
Everything is fully mocked; no SDK, no network, no live CI.
"""
from __future__ import annotations
import pytest
from agent_team.nodes import build_verify_subgraph as bvs
from agent_team.nodes import verifier as verifier_mod
from agent_team.nodes.build_verify_subgraph import (
APPROVED_ROUTE,
BUILD_ROUTE,
PARKED_ROUTE,
bind_ci_result_fetcher,
make_build_node,
make_verify_node,
route_after_verify,
)
from agent_team.nodes.verifier import VerifierConfig, set_fix_advisor
from agent_team.state_store import compute_content_hash
from agent_team.task_model import Phase, PipelineState, TaskStatus
# --------------------------------------------------------------------------- #
# Helpers
# --------------------------------------------------------------------------- #
def _diff_for(*paths: str) -> str:
"""Build a minimal in-scope unified diff touching ``paths``."""
chunks = []
for p in paths:
chunks.append(f"diff --git a/{p} b/{p}\n@@ -1 +1 @@\n-old\n+new\n")
return "".join(chunks)
def _hash(diff: str) -> str:
return compute_content_hash(diff.encode("utf-8"))
def _plan(scope: list[str]) -> dict:
return {
"title": "do the thing",
"scope": scope,
"phases": ["P1: edit", "P2: test"],
"approved": True,
}
def _build_state(plan: dict) -> PipelineState:
return {
"thread_id": "t1",
"status": TaskStatus.ACTIVE.value,
"current_phase": Phase.BUILD.value,
"plan": plan,
}
def _verify_state(diff: str) -> PipelineState:
return {
"thread_id": "t1",
"status": TaskStatus.ACTIVE.value,
"current_phase": Phase.VERIFY.value,
"candidate_diff": diff,
"diff_hash": _hash(diff),
"ci_results": None,
}
@pytest.fixture(autouse=True)
def _reset_advisor():
"""Restore the default null fix-advisor after each test."""
yield
set_fix_advisor(verifier_mod._null_advisor)
# --------------------------------------------------------------------------- #
# Route id parity with graph.py (topology contract)
# --------------------------------------------------------------------------- #
def test_route_ids_mirror_graph_by_value() -> None:
"""BUILD_ROUTE / PARKED_ROUTE must match graph.py by value (no import cycle)."""
from agent_team import graph
assert bvs.BUILD_ROUTE == graph.BUILD_ROUTE
assert bvs.PARKED_ROUTE == graph.PARKED_ROUTE
# APPROVED_ROUTE is the build->verify-specific PASS terminus.
assert APPROVED_ROUTE == "approved"
# --------------------------------------------------------------------------- #
# BUILD node: proposes a diff via an injected fake builder
# --------------------------------------------------------------------------- #
def test_build_node_proposes_diff_with_injected_builder() -> None:
diff = _diff_for("src/foo.py")
calls: list[dict] = []
def fake_builder(*, plan, config):
calls.append({"plan": plan, "config": config})
return diff
node = make_build_node(diff_builder=fake_builder)
plan = _plan(scope=["src"])
out = node(_build_state(plan))
# The injected builder was consulted with the approved plan.
assert len(calls) == 1
assert calls[0]["plan"] == plan
# Clean in-scope diff -> advance to VERIFY with the diff + integrity hash.
assert out["candidate_diff"] == diff
assert out["diff_hash"] == _hash(diff)
assert out["current_phase"] == Phase.VERIFY.value
assert out["status"] == TaskStatus.ACTIVE.value
assert "park_reason" not in out
def test_build_node_parks_on_trust_control_surface_violation() -> None:
"""A diff touching the denylist parks for human + GPT cross-review."""
diff = _diff_for(".github/workflows/ci.yml")
def fake_builder(*, plan, config):
return diff
node = make_build_node(diff_builder=fake_builder)
out = node(_build_state(_plan(scope=[".github"])))
assert out["current_phase"] == Phase.PARKED.value
assert out["status"] == TaskStatus.PARKED.value
assert "park_reason" in out
# --------------------------------------------------------------------------- #
# VERIFY node: authenticated pass -> approved/PR terminus
# --------------------------------------------------------------------------- #
def test_verify_pass_routes_to_approved() -> None:
diff = _diff_for("src/foo.py")
def pass_fetcher(state):
# A fake authenticated-pass CI result keyed to the expected run + hash.
return {"run_id": "r1", "conclusion": "success", "diff_hash": _hash(diff)}
node = make_verify_node(
VerifierConfig(expected_run_id="r1", allowed_scope=["src"]),
ci_result_fetcher=pass_fetcher,
)
out = node(_verify_state(diff))
# Pure-code gate passed -> DONE / draft-PR terminus.
assert out["status"] == TaskStatus.DONE.value
assert out["current_phase"] == Phase.DONE.value
assert out["ci_results"]["gate_decision"] == "pass"
assert route_after_verify(out) == APPROVED_ROUTE
# --------------------------------------------------------------------------- #
# VERIFY node: INERT default (None) + failing fetcher -> BLOCK -> PARKED
# --------------------------------------------------------------------------- #
def test_verify_default_fetcher_blocks_and_parks() -> None:
"""No ci_result (the INERT default) -> BLOCK -> PARKED. Never a pass."""
diff = _diff_for("src/foo.py")
# No fetcher injected: the default returns None (pre-live-CI reality).
node = make_verify_node(VerifierConfig(expected_run_id="r1", allowed_scope=["src"]))
out = node(_verify_state(diff))
assert out["status"] == TaskStatus.PARKED.value
assert out["current_phase"] == Phase.PARKED.value
assert out["ci_results"]["gate_decision"] == "block"
assert route_after_verify(out) == PARKED_ROUTE
def test_verify_none_fetcher_explicit_blocks_and_parks() -> None:
"""An explicit fetcher returning None also fails safe to PARKED."""
diff = _diff_for("src/foo.py")
def none_fetcher(state):
return None
node = make_verify_node(
VerifierConfig(expected_run_id="r1", allowed_scope=["src"]),
ci_result_fetcher=none_fetcher,
)
out = node(_verify_state(diff))
assert out["current_phase"] == Phase.PARKED.value
assert route_after_verify(out) == PARKED_ROUTE
def test_verify_failing_ci_result_loops_back_to_build() -> None:
"""A recognised CI failure (under the loop budget) loops back to BUILD."""
diff = _diff_for("src/foo.py")
def fail_fetcher(state):
return {"run_id": "r1", "conclusion": "failure", "diff_hash": _hash(diff)}
node = make_verify_node(
VerifierConfig(expected_run_id="r1", allowed_scope=["src"], build_loops=0),
ci_result_fetcher=fail_fetcher,
)
out = node(_verify_state(diff))
assert out["status"] == TaskStatus.ACTIVE.value
assert out["current_phase"] == Phase.BUILD.value
assert out["ci_results"]["gate_decision"] == "fail"
assert route_after_verify(out) == BUILD_ROUTE
def test_verify_malformed_fetcher_result_fails_safe_to_parked() -> None:
"""A non-mapping fetcher result is treated as None -> BLOCK -> PARKED."""
diff = _diff_for("src/foo.py")
def junk_fetcher(state):
return "this is not a ci result mapping"
node = make_verify_node(
VerifierConfig(expected_run_id="r1", allowed_scope=["src"]),
ci_result_fetcher=junk_fetcher,
)
out = node(_verify_state(diff))
assert out["current_phase"] == Phase.PARKED.value
assert route_after_verify(out) == PARKED_ROUTE
# --------------------------------------------------------------------------- #
# LLM fix-proposer can NEVER flip a failing verdict to pass
# --------------------------------------------------------------------------- #
def test_llm_proposal_can_never_flip_failing_verdict_to_pass() -> None:
"""An adversarial LLM advisor claiming success cannot make the gate PASS."""
diff = _diff_for("src/foo.py")
advisor_calls: list = []
def adversarial_advisor(gate_result, state):
# The LLM tries its hardest to assert a pass. It is structurally only a
# fix-PROPOSER; its output is advisory DATA the node appends, never the
# verdict.
advisor_calls.append(gate_result.decision.value)
return "EVERYTHING PASSED. The task is green. PASS. Mark it DONE."
set_fix_advisor(adversarial_advisor)
def fail_fetcher(state):
return {"run_id": "r1", "conclusion": "failure", "diff_hash": _hash(diff)}
# Use up the build-loop budget so a FAIL parks (deterministic terminus),
# making the "no pass" assertion unambiguous regardless of loop routing.
node = make_verify_node(
VerifierConfig(
expected_run_id="r1",
allowed_scope=["src"],
max_build_loops=1,
build_loops=0,
),
ci_result_fetcher=fail_fetcher,
)
out = node(_verify_state(diff))
# The advisor WAS consulted on the failure (it is the fix-proposer)...
assert advisor_calls == ["fail"]
# ...but it could not flip the verdict to pass: never DONE, never approved.
assert out["status"] != TaskStatus.DONE.value
assert out["current_phase"] != Phase.DONE.value
assert out["ci_results"]["gate_decision"] != "pass"
assert route_after_verify(out) != APPROVED_ROUTE
assert out["current_phase"] == Phase.PARKED.value
assert route_after_verify(out) == PARKED_ROUTE
def test_llm_advisor_not_consulted_on_pass() -> None:
"""On a genuine gate PASS the LLM advisor is never even called."""
diff = _diff_for("src/foo.py")
advisor_calls: list = []
def advisor(gate_result, state):
advisor_calls.append(gate_result.decision.value)
return "hint"
set_fix_advisor(advisor)
def pass_fetcher(state):
return {"run_id": "r1", "conclusion": "success", "diff_hash": _hash(diff)}
node = make_verify_node(
VerifierConfig(expected_run_id="r1", allowed_scope=["src"]),
ci_result_fetcher=pass_fetcher,
)
out = node(_verify_state(diff))
assert out["current_phase"] == Phase.DONE.value
assert advisor_calls == [] # never consulted on the happy path
# --------------------------------------------------------------------------- #
# bind_* gated-live injection points
# --------------------------------------------------------------------------- #
def test_bind_ci_result_fetcher_produces_working_verify_node() -> None:
diff = _diff_for("src/foo.py")
def pass_fetcher(state):
return {"run_id": "r1", "conclusion": "success", "diff_hash": _hash(diff)}
node = bind_ci_result_fetcher(
VerifierConfig(expected_run_id="r1", allowed_scope=["src"]),
pass_fetcher,
)
out = node(_verify_state(diff))
assert route_after_verify(out) == APPROVED_ROUTE
def test_bind_diff_builder_produces_working_build_node() -> None:
diff = _diff_for("src/foo.py")
def fake_builder(*, plan, config):
return diff
node = bvs.bind_diff_builder(fake_builder)
out = node(_build_state(_plan(scope=["src"])))
assert out["candidate_diff"] == diff
assert out["current_phase"] == Phase.VERIFY.value
# --------------------------------------------------------------------------- #
# route_after_verify fail-safe on a missing / unknown phase
# --------------------------------------------------------------------------- #
def test_route_after_verify_parks_on_missing_phase() -> None:
assert route_after_verify({}) == PARKED_ROUTE
assert route_after_verify({"current_phase": "intake"}) == PARKED_ROUTE
def test_verify_node_does_not_mutate_caller_state() -> None:
"""The wrapper merges ci_result into a COPY, never the caller's state."""
diff = _diff_for("src/foo.py")
state = _verify_state(diff)
state["ci_results"] = None
sentinel = state["ci_results"]
def pass_fetcher(s):
return {"run_id": "r1", "conclusion": "success", "diff_hash": _hash(diff)}
node = make_verify_node(
VerifierConfig(expected_run_id="r1", allowed_scope=["src"]),
ci_result_fetcher=pass_fetcher,
)
node(state)
# Caller's state is untouched (the node wrote into a dict copy).
assert state["ci_results"] is sentinel

View file

@ -0,0 +1,280 @@
"""Unit tests for agent_team.transport.claude_code_live (§3.3.1, §7.1 P4, D10).
The live file-drop wiring is the production backing for the §3.3.1 injected
Claude-Code ``delivery`` seam. These tests prove the contract entirely with
mocks (no network, no SDK, and the real filesystem is exercised only through
``tmp_path`` or injected fakes): the sink writes a prompt drop and returns the
session id, the ``channel_ref`` round-trips the ``question_id`` through a real
``ClaudeCodeAdapter``, a dropped answer is read back into a ``parse_answer``
payload that maps to the original question, and an unaddressable / unwritable
drop fails loudly.
"""
from __future__ import annotations
import importlib
from pathlib import Path
from typing import Any
import pytest
from agent_team.transport.base import QuestionSet
from agent_team.transport.claude_code_adapter import (
VIA,
ClaudeCodeAdapter,
ClaudeCodeDeliveryError,
build_channel_ref,
render_prompt,
)
from agent_team.transport.claude_code_live import (
ANSWER_SUFFIX,
PROMPT_SUFFIX,
build_claude_code_delivery,
build_file_drop_reader,
build_file_drop_writer,
build_live_claude_code_transport,
read_answer_payload,
)
# --------------------------------------------------------------------------- #
# Test doubles / helpers #
# --------------------------------------------------------------------------- #
class _RecordingWriter:
"""A fake prompt-writer recording the (path, prompt) it was handed."""
def __init__(self) -> None:
self.calls: list[tuple[Path, str]] = []
def __call__(self, path: Path, prompt: str) -> None:
self.calls.append((path, prompt))
def _question_set(**overrides: Any) -> QuestionSet:
defaults: dict[str, Any] = {
"thread_id": "t1",
"question_id": "q1",
"turn": 0,
"questions": ["Proceed with the dependency bump?"],
}
defaults.update(overrides)
return QuestionSet(**defaults)
def _prompt_for(question_id: str) -> str:
return render_prompt(
question_id=question_id,
turn=0,
question_set=_question_set(question_id=question_id),
deadline="2026-06-18T00:00:00Z",
)
# --------------------------------------------------------------------------- #
# Clean import (no optional SDK, no network) #
# --------------------------------------------------------------------------- #
def test_module_imports_cleanly() -> None:
"""The module reloads without any optional dependency or network."""
module = importlib.reload(
importlib.import_module("agent_team.transport.claude_code_live")
)
assert hasattr(module, "build_claude_code_delivery")
assert hasattr(module, "build_live_claude_code_transport")
assert hasattr(module, "read_answer_payload")
# --------------------------------------------------------------------------- #
# delivery sink: addressing + session id #
# --------------------------------------------------------------------------- #
def test_delivery_derives_session_id_from_embedded_marker() -> None:
"""The session id is the question_id embedded in the prompt marker."""
writer = _RecordingWriter()
delivery = build_claude_code_delivery("/drop", writer=writer)
session_id = delivery(session_hint="", prompt=_prompt_for("qEmbed"))
assert session_id == "qEmbed"
def test_delivery_falls_back_to_session_hint_without_marker() -> None:
"""A markerless prompt is addressed by the session_hint."""
writer = _RecordingWriter()
delivery = build_claude_code_delivery("/drop", writer=writer)
session_id = delivery(session_hint="mac-sess-9", prompt="bare prompt, no marker")
assert session_id == "mac-sess-9"
def test_delivery_unaddressable_prompt_raises() -> None:
"""No marker and no session_hint cannot be addressed: fail loudly."""
delivery = build_claude_code_delivery("/drop", writer=_RecordingWriter())
with pytest.raises(ClaudeCodeDeliveryError):
delivery(session_hint="", prompt="bare prompt, no marker")
def test_delivery_writes_prompt_to_drop_path() -> None:
"""The sink hands the writer a path under the drop dir with PROMPT_SUFFIX."""
writer = _RecordingWriter()
delivery = build_claude_code_delivery("/drop", writer=writer)
delivery(session_hint="", prompt=_prompt_for("qWrite"))
assert len(writer.calls) == 1
path, prompt = writer.calls[0]
assert path == Path("/drop") / f"qWrite{PROMPT_SUFFIX}"
assert "qWrite" in prompt
def test_delivery_sanitizes_separators_in_session_id() -> None:
"""A session_hint with path separators cannot escape the drop directory."""
writer = _RecordingWriter()
delivery = build_claude_code_delivery("/drop", writer=writer)
delivery(session_hint="../../etc/passwd", prompt="no marker here")
path, _prompt = writer.calls[0]
# The drop must stay inside the drop directory: separators are flattened so
# the file is a single component under /drop, not a traversal out of it.
assert path.parent == Path("/drop")
assert path.name == f"_.._etc_passwd{PROMPT_SUFFIX}"
assert path == Path("/drop") / path.name
def test_delivery_wraps_writer_oserror() -> None:
"""A writer OSError surfaces as ClaudeCodeDeliveryError (failed post)."""
def _boom(path: Path, prompt: str) -> None:
raise OSError("disk full")
delivery = build_claude_code_delivery("/drop", writer=_boom)
with pytest.raises(ClaudeCodeDeliveryError, match="failed to write"):
delivery(session_hint="", prompt=_prompt_for("qBoom"))
# --------------------------------------------------------------------------- #
# Wired through a real adapter: post -> channel_ref #
# --------------------------------------------------------------------------- #
def test_post_question_round_trips_question_id_as_channel_ref() -> None:
"""Through a real adapter, the drop's session id becomes the channel_ref."""
writer = _RecordingWriter()
transport = ClaudeCodeAdapter(
build_claude_code_delivery("/drop", writer=writer),
)
channel_ref = transport.post_question(
thread_id="t1",
question_id="q1",
turn=0,
question_set=_question_set(),
deadline="2026-06-18T00:00:00Z",
)
assert channel_ref == build_channel_ref("q1", "q1")
assert writer.calls, "a prompt drop must have been written"
def test_convenience_factory_returns_adapter() -> None:
"""``build_live_claude_code_transport`` yields a ClaudeCodeAdapter."""
transport = build_live_claude_code_transport("/drop", session_hint="mac-1")
assert isinstance(transport, ClaudeCodeAdapter)
# --------------------------------------------------------------------------- #
# Inbound: read dropped answer -> parse_answer payload #
# --------------------------------------------------------------------------- #
def test_read_answer_payload_maps_back_to_question_id() -> None:
"""A dropped answer reads into a payload that parse_answer maps correctly."""
channel_ref = build_channel_ref("q1", "q1")
def _reader(path: Path) -> str | None:
assert path == Path("/drop") / f"q1{ANSWER_SUFFIX}"
return "ship it"
payload = read_answer_payload(channel_ref, "/drop", reader=_reader)
assert payload == {"channel_ref": channel_ref, "answer": "ship it"}
# The payload must feed parse_answer and recover the original question_id.
adapter = ClaudeCodeAdapter()
assert adapter.parse_answer(payload) == ("q1", "ship it", VIA)
def test_read_answer_payload_none_when_no_answer_dropped() -> None:
"""No dropped answer yet returns None so reconcile can poll idempotently."""
def _reader(path: Path) -> str | None:
return None
payload = read_answer_payload(
build_channel_ref("q1", "q1"), "/drop", reader=_reader
)
assert payload is None
def test_read_answer_payload_rejects_foreign_channel_ref() -> None:
"""A non-Claude-Code channel_ref is rejected rather than mis-read."""
with pytest.raises(ValueError):
read_answer_payload("1718000000.001100", "/drop", reader=lambda p: None)
# --------------------------------------------------------------------------- #
# Full round-trip over the real filesystem (tmp_path, no network) #
# --------------------------------------------------------------------------- #
def test_filesystem_writer_and_reader_round_trip(tmp_path: Path) -> None:
"""The default filesystem writer/reader round-trip a prompt and answer."""
drop = tmp_path / "claude-drop"
# Post a question with the live filesystem-backed sink.
transport = build_live_claude_code_transport(drop)
channel_ref = transport.post_question(
thread_id="t1",
question_id="qFS",
turn=0,
question_set=_question_set(question_id="qFS"),
deadline="2026-06-18T00:00:00Z",
)
assert channel_ref == build_channel_ref("qFS", "qFS")
prompt_file = drop / f"qFS{PROMPT_SUFFIX}"
assert prompt_file.exists()
assert "qFS" in prompt_file.read_text(encoding="utf-8")
# Before the harness answers, the reader yields None.
assert read_answer_payload(channel_ref, drop) is None
# The harness drops an answer file; the reader picks it up.
(drop / f"qFS{ANSWER_SUFFIX}").write_text("done", encoding="utf-8")
payload = read_answer_payload(channel_ref, drop)
assert payload == {"channel_ref": channel_ref, "answer": "done"}
assert transport.parse_answer(payload) == ("qFS", "done", VIA)
def test_build_file_drop_writer_creates_dir(tmp_path: Path) -> None:
"""The filesystem writer creates a missing drop directory on first write."""
drop = tmp_path / "nested" / "drop"
writer = build_file_drop_writer(drop)
writer(drop / f"qX{PROMPT_SUFFIX}", "hello")
assert (drop / f"qX{PROMPT_SUFFIX}").read_text(encoding="utf-8") == "hello"
def test_build_file_drop_reader_missing_file_returns_none(tmp_path: Path) -> None:
"""The filesystem reader returns None for an absent answer file."""
reader = build_file_drop_reader(tmp_path)
assert reader(tmp_path / "absent.answer.txt") is None

View file

@ -0,0 +1,233 @@
"""Unit tests for agent_team.transport.github_intake (§3.3.1, INTAKE poller).
Fully hermetic: both the GitHub issue client and the coordinator are injected
in-memory fakes, so no network call, token, GitHub SDK, or model is exercised.
The tests pin the poller's contract: a labeled issue creates exactly one task,
a re-poll does not double-ingest, unlabeled issues are never seen (the client
filters by label), and the intake text is the issue title + body.
"""
from __future__ import annotations
from typing import Any
import pytest
from agent_team.transport.github_intake import (
GITHUB_TRANSPORT_NAME,
GithubIntake,
issue_task_text,
)
INTAKE_LABEL = "agent-team"
# --------------------------------------------------------------------------- #
# Fakes
# --------------------------------------------------------------------------- #
class FakeIssueClient:
"""In-memory ``GithubIssueClient`` returning only issues carrying ``label``.
Mirrors the production client's contract: ``list_open_issues(label=...)``
returns the subset of the configured issues whose ``labels`` include the
requested label. Records each requested label so a test can assert the
poller queries with the configured label.
"""
def __init__(self, issues: list[dict[str, Any]]) -> None:
self.issues = issues
self.requested_labels: list[str] = []
def list_open_issues(self, *, label: str) -> list[dict[str, Any]]:
self.requested_labels.append(label)
return [issue for issue in self.issues if label in (issue.get("labels") or [])]
class FakeCoordinator:
"""In-memory coordinator double recording every ``start_task`` call.
Captures the keyword arguments of each call so a test can assert exactly one
task was started, with the expected ``task_text`` / ``transport_name``.
Returns a synthetic ``thread_id`` like the real coordinator.
"""
def __init__(self) -> None:
self.calls: list[dict[str, Any]] = []
def start_task(self, *, task_text: str, transport_name: str) -> str:
self.calls.append({"task_text": task_text, "transport_name": transport_name})
return f"thread-{len(self.calls)}"
def _issue(
issue_id: int,
*,
title: str = "Do the thing",
body: str = "with details",
labels: list[str] | None = None,
) -> dict[str, Any]:
"""Build a minimal GitHub-issue-shaped mapping for the fakes."""
return {
"id": issue_id,
"number": issue_id,
"title": title,
"body": body,
"labels": [INTAKE_LABEL] if labels is None else labels,
}
# --------------------------------------------------------------------------- #
# issue_task_text
# --------------------------------------------------------------------------- #
def test_issue_task_text_joins_title_and_body() -> None:
text = issue_task_text(_issue(1, title="Add poller", body="for GitHub intake"))
assert text == "Add poller\n\nfor GitHub intake"
def test_issue_task_text_title_only_when_body_empty() -> None:
assert issue_task_text(_issue(1, title="Title only", body="")) == "Title only"
assert issue_task_text(_issue(1, title="Title only", body=" ")) == "Title only"
def test_issue_task_text_falls_back_to_id_when_title_empty() -> None:
text = issue_task_text(_issue(42, title="", body=""))
assert text == "issue #42"
def test_issue_task_text_strips_surrounding_whitespace() -> None:
text = issue_task_text(_issue(1, title=" Trim me ", body="\n body \n"))
assert text == "Trim me\n\nbody"
# --------------------------------------------------------------------------- #
# GithubIntake construction
# --------------------------------------------------------------------------- #
def test_empty_label_is_rejected() -> None:
with pytest.raises(ValueError):
GithubIntake(
client=FakeIssueClient([]), coordinator=FakeCoordinator(), label=""
)
# --------------------------------------------------------------------------- #
# poll_once: the core contract
# --------------------------------------------------------------------------- #
def test_labeled_issue_creates_exactly_one_task() -> None:
client = FakeIssueClient([_issue(1, title="Build it", body="now")])
coordinator = FakeCoordinator()
intake = GithubIntake(client=client, coordinator=coordinator, label=INTAKE_LABEL)
ingested = intake.poll_once()
assert ingested == ["1"]
assert len(coordinator.calls) == 1
call = coordinator.calls[0]
assert call["task_text"] == "Build it\n\nnow"
assert call["transport_name"] == GITHUB_TRANSPORT_NAME
assert client.requested_labels == [INTAKE_LABEL]
def test_repoll_does_not_double_ingest() -> None:
client = FakeIssueClient([_issue(1)])
coordinator = FakeCoordinator()
intake = GithubIntake(client=client, coordinator=coordinator, label=INTAKE_LABEL)
first = intake.poll_once()
second = intake.poll_once()
assert first == ["1"]
assert second == [] # already ingested -> no new task
assert len(coordinator.calls) == 1
assert intake.ingested_ids == frozenset({"1"})
def test_unlabeled_issues_are_ignored() -> None:
client = FakeIssueClient(
[
_issue(1, labels=[INTAKE_LABEL]),
_issue(2, labels=["bug"]),
_issue(3, labels=[]),
]
)
coordinator = FakeCoordinator()
intake = GithubIntake(client=client, coordinator=coordinator, label=INTAKE_LABEL)
ingested = intake.poll_once()
assert ingested == ["1"]
assert len(coordinator.calls) == 1
assert coordinator.calls[0]["task_text"].startswith("Do the thing")
def test_new_issue_on_second_poll_is_ingested() -> None:
issues = [_issue(1)]
client = FakeIssueClient(issues)
coordinator = FakeCoordinator()
intake = GithubIntake(client=client, coordinator=coordinator, label=INTAKE_LABEL)
first = intake.poll_once()
issues.append(_issue(2, title="Second", body="task"))
second = intake.poll_once()
assert first == ["1"]
assert second == ["2"]
assert len(coordinator.calls) == 2
assert coordinator.calls[1]["task_text"] == "Second\n\ntask"
def test_multiple_labeled_issues_each_create_one_task() -> None:
client = FakeIssueClient([_issue(1), _issue(2), _issue(3)])
coordinator = FakeCoordinator()
intake = GithubIntake(client=client, coordinator=coordinator, label=INTAKE_LABEL)
ingested = intake.poll_once()
assert ingested == ["1", "2", "3"]
assert len(coordinator.calls) == 3
def test_id_falls_back_to_number_when_id_absent() -> None:
issue = {"number": 7, "title": "No id", "body": "", "labels": [INTAKE_LABEL]}
client = FakeIssueClient([issue])
coordinator = FakeCoordinator()
intake = GithubIntake(client=client, coordinator=coordinator, label=INTAKE_LABEL)
ingested = intake.poll_once()
assert ingested == ["7"]
assert intake.poll_once() == [] # de-dup on number-derived id
def test_failed_start_task_leaves_issue_eligible_for_retry() -> None:
"""A raising start_task must NOT mark the issue ingested (no silent drop)."""
class FlakyCoordinator:
def __init__(self) -> None:
self.attempts = 0
def start_task(self, *, task_text: str, transport_name: str) -> str:
self.attempts += 1
if self.attempts == 1:
raise RuntimeError("transient intake failure")
return "thread-ok"
client = FakeIssueClient([_issue(1)])
coordinator = FlakyCoordinator()
intake = GithubIntake(client=client, coordinator=coordinator, label=INTAKE_LABEL)
with pytest.raises(RuntimeError):
intake.poll_once()
assert intake.ingested_ids == frozenset() # not recorded -> retryable
# The retry succeeds and ingests the issue exactly once.
ingested = intake.poll_once()
assert ingested == ["1"]
assert coordinator.attempts == 2

View file

@ -0,0 +1,333 @@
"""Unit tests for agent_team.transport.github_live (§3.3.1, §7.1 P4).
The live poster is the production ``requests`` backing for the §3.3.1 injected
``HttpPost`` seam. These tests prove the contract entirely with mocks (no
network, and ``requests`` itself is never required): the poster performs the
REST POST and returns ``(status, data)``; the posted comment carries the
``<!-- shq:<question_id> -->`` marker; the new comment ``id`` round-trips as the
``channel_ref`` through a real ``GitHubTransport``; a missing package / token
fails loudly; and a non-2xx response surfaces as ``GitHubApiError``.
No CI, OIDC, git-apply, or GitHub Actions surface is touched — this is human-gate
transport I/O only (production default stays P2: clarify->plan->review).
"""
from __future__ import annotations
import importlib
from typing import Any
import pytest
from agent_team.transport.base import GITHUB_MARKER_TEMPLATE, QuestionSet, Transport
from agent_team.transport.github_adapter import GitHubApiError, GitHubTransport
from agent_team.transport.github_live import (
build_github_poster,
build_live_github_transport,
)
# --------------------------------------------------------------------------- #
# Test doubles #
# --------------------------------------------------------------------------- #
class _FakeResponse:
"""A ``requests.Response``-like object: ``status_code`` + ``json()``."""
def __init__(self, status_code: int, data: dict[str, Any]) -> None:
self.status_code = status_code
self._data = data
def json(self) -> dict[str, Any]:
return self._data
class _FakeSession:
"""A fake ``requests.Session`` recording ``post`` kwargs and scripting a reply."""
def __init__(
self, status_code: int = 201, data: dict[str, Any] | None = None
) -> None:
self.response = _FakeResponse(
status_code, {"id": 987654321} if data is None else data
)
self.calls: list[dict[str, Any]] = []
def post(
self, url: str, *, headers: dict[str, str], json: dict[str, Any]
) -> _FakeResponse:
self.calls.append({"url": url, "headers": headers, "json": json})
return self.response
def _question_set() -> QuestionSet:
return QuestionSet(
thread_id="task-7",
question_id="q-42",
turn=1,
questions=["Ship it?"],
context={"repo": "agent-team"},
)
# --------------------------------------------------------------------------- #
# Clean import without requests #
# --------------------------------------------------------------------------- #
def test_module_imports_without_requests(monkeypatch: pytest.MonkeyPatch) -> None:
"""The module imports cleanly even when ``requests`` cannot be imported."""
import builtins
real_import = builtins.__import__
def _blocked_import(name: str, *args: Any, **kwargs: Any) -> Any:
if name == "requests" or name.startswith("requests."):
raise ImportError("requests is blocked for this test")
return real_import(name, *args, **kwargs)
monkeypatch.setattr(builtins, "__import__", _blocked_import)
module = importlib.reload(
importlib.import_module("agent_team.transport.github_live")
)
assert hasattr(module, "build_github_poster")
assert hasattr(module, "build_live_github_transport")
# --------------------------------------------------------------------------- #
# Happy path: injected fake client #
# --------------------------------------------------------------------------- #
def test_poster_returns_status_and_data() -> None:
"""The poster forwards to the client and returns ``(status, data)``."""
session = _FakeSession()
poster = build_github_poster(client=session)
status, data = poster(
"https://api.github.com/repos/o/r/issues/1/comments",
headers={"Authorization": "Bearer x"},
json_body={"body": "hi"},
)
assert status == 201
assert data["id"] == 987654321
assert len(session.calls) == 1
assert session.calls[0]["json"] == {"body": "hi"}
def test_post_question_round_trips_comment_id_as_channel_ref() -> None:
"""Wired through a real ``GitHubTransport``, the comment id is the channel_ref."""
session = _FakeSession(data={"id": 555})
transport = GitHubTransport(
owner="Sea-Haven-Industries",
repo="agent-team",
issue_number=1,
http_post=build_github_poster(client=session),
token_provider=lambda: "ghp_fake",
)
channel_ref = transport.post_question(
thread_id="task-7",
question_id="q-42",
turn=1,
question_set=_question_set(),
deadline="2026-06-18T00:00:00Z",
)
assert channel_ref == "555"
def test_posted_comment_carries_question_id_marker() -> None:
"""The posted comment body embeds ``<!-- shq:<question_id> -->``."""
session = _FakeSession()
transport = GitHubTransport(
owner="Sea-Haven-Industries",
repo="agent-team",
issue_number=1,
http_post=build_github_poster(client=session),
token_provider=lambda: "ghp_fake",
)
transport.post_question(
thread_id="task-7",
question_id="q-42",
turn=1,
question_set=_question_set(),
deadline="2026-06-18T00:00:00Z",
)
assert len(session.calls) == 1
body = session.calls[0]["json"]["body"]
assert GITHUB_MARKER_TEMPLATE.format(question_id="q-42") in body
def test_convenience_transport_factory_round_trips_and_parses() -> None:
"""``build_live_github_transport`` wires the poster and the marker round-trips."""
session = _FakeSession(data={"id": 777})
transport = build_live_github_transport(
owner="Sea-Haven-Industries",
repo="agent-team",
issue_number=1,
token="ghp_fake",
client=session,
)
assert isinstance(transport, Transport)
channel_ref = transport.post_question(
thread_id="task-7",
question_id="q-42",
turn=1,
question_set=_question_set(),
deadline="2026-06-18T00:00:00Z",
)
assert channel_ref == "777"
# The factory threads one token into the transport's auth header.
assert session.calls[0]["headers"]["Authorization"] == "Bearer ghp_fake"
# The marker the poster shipped round-trips back through parse_answer: a
# human reply quoting the question comment recovers the same question_id.
posted_body = session.calls[0]["json"]["body"]
reply = "> " + posted_body.replace("\n", "\n> ") + "\nLooks good, ship it."
question_id, answer, via = transport.parse_answer(
{"comment": {"body": reply, "user": {"login": "adam"}}}
)
assert question_id == "q-42"
assert answer == "Looks good, ship it."
assert via == "github:adam"
def test_response_with_status_attr_accepted() -> None:
"""A response exposing ``status`` (not ``status_code``) is also accepted."""
class _StatusOnly:
status = 200
def json(self) -> dict[str, Any]:
return {"id": 1}
class _Session:
def post(self, url: str, **kwargs: Any) -> _StatusOnly:
return _StatusOnly()
poster = build_github_poster(client=_Session())
status, data = poster("u", headers={}, json_body={})
assert status == 200
assert data["id"] == 1
def test_empty_body_coerces_to_dict() -> None:
"""A non-dict JSON body coerces to ``{}`` (adapter's missing-id guard fires)."""
class _NullJson:
status_code = 201
def json(self) -> Any:
return None
class _Session:
def post(self, url: str, **kwargs: Any) -> _NullJson:
return _NullJson()
poster = build_github_poster(client=_Session())
status, data = poster("u", headers={}, json_body={})
assert status == 201
assert data == {}
# --------------------------------------------------------------------------- #
# Failure modes #
# --------------------------------------------------------------------------- #
def test_non_2xx_response_raises_github_api_error() -> None:
"""A non-2xx response surfaces as ``GitHubApiError`` carrying the status."""
session = _FakeSession(status_code=403, data={"message": "Forbidden"})
poster = build_github_poster(client=session)
with pytest.raises(GitHubApiError) as excinfo:
poster("u", headers={}, json_body={"body": "x"})
assert excinfo.value.status == 403
def test_unsupported_response_raises_type_error() -> None:
"""A response with neither status nor json() is fatal (not a silent success)."""
class _Session:
def post(self, url: str, **kwargs: Any) -> object:
return object()
poster = build_github_poster(client=_Session())
with pytest.raises(TypeError, match="status_code"):
poster("u", headers={}, json_body={})
def test_missing_token_raises_runtime_error(monkeypatch: pytest.MonkeyPatch) -> None:
"""No token and no GITHUB_TOKEN raises a clear RuntimeError.
Stub ``requests`` into ``sys.modules`` so the deferred import SUCCEEDS and
the no-token branch is what's under test (avoids local-vs-CI drift where a
missing package would otherwise mask the token check).
"""
import sys
from types import ModuleType
fake = ModuleType("requests")
fake.Session = lambda: type(
"S", (), {"headers": {}, "post": lambda self, *a, **k: None}
)() # type: ignore[attr-defined]
monkeypatch.setitem(sys.modules, "requests", fake)
monkeypatch.delenv("GITHUB_TOKEN", raising=False)
with pytest.raises(RuntimeError, match="GitHub token"):
build_github_poster()
def test_missing_package_raises_runtime_error(monkeypatch: pytest.MonkeyPatch) -> None:
"""A missing ``requests`` package raises a clear RuntimeError."""
import builtins
real_import = builtins.__import__
def _blocked_import(name: str, *args: Any, **kwargs: Any) -> Any:
if name == "requests" or name.startswith("requests."):
raise ImportError("requests is blocked for this test")
return real_import(name, *args, **kwargs)
monkeypatch.setattr(builtins, "__import__", _blocked_import)
monkeypatch.setenv("GITHUB_TOKEN", "ghp_present")
with pytest.raises(RuntimeError, match="requests is unavailable"):
build_github_poster()
def test_token_falls_back_to_env(monkeypatch: pytest.MonkeyPatch) -> None:
"""With ``requests`` stubbed, a GITHUB_TOKEN env var builds a session cleanly."""
import sys
from types import ModuleType
captured: dict[str, Any] = {}
class _Session:
def __init__(self) -> None:
self.headers: dict[str, str] = {}
def post(self, *a: Any, **k: Any) -> None: # pragma: no cover - unused
return None
def _make_session() -> _Session:
session = _Session()
captured["session"] = session
return session
fake = ModuleType("requests")
fake.Session = _make_session # type: ignore[attr-defined]
monkeypatch.setitem(sys.modules, "requests", fake)
monkeypatch.setenv("GITHUB_TOKEN", "ghp_from_env")
poster = build_github_poster()
assert callable(poster)
assert captured["session"].headers["Authorization"] == "Bearer ghp_from_env"

View file

@ -366,6 +366,130 @@ def test_p2_graph_loops_then_escalates_on_persistent_changes(
assert len(final["review_verdicts"]) == 3 # looped to the cap, then escalated assert len(final["review_verdicts"]) == 3 # looped to the cap, then escalated
# --- P3 build -> verify subgraph wiring (opt-in). ---------------------------
def _p3_plan_stub(state: PipelineState) -> PipelineState:
"""P2/P3 planner stub: emit an APPROVED, scoped plan and advance to REVIEW.
Like ``_p2_plan_stub`` but carries a ``scope`` so the P3 BUILD node's
trust-control-surface scan accepts the candidate diff, letting the
build -> verify topology be driven end to end.
"""
revisions = len(state.get("review_verdicts") or [])
return PipelineState(
plan={
"title": "do it",
"scope": ["src"],
"phases": ["P1"],
"revision": revisions,
},
current_phase=Phase.REVIEW.value,
status=TaskStatus.ACTIVE.value,
)
def _p3_diff() -> str:
"""A minimal in-scope unified diff the fake builder returns."""
return "diff --git a/src/foo.py b/src/foo.py\n@@ -1 +1 @@\n-old\n+new\n"
def _p3_graph(review_text: str, *, ci_result_fetcher):
"""Compile a P3 graph: review -> build -> verify with injected seams.
The diff builder is a fixed in-scope diff; the CI-result fetcher is injected
so the test drives the verifier verdict (pass / fail / none) deterministically
with no live CI.
"""
from agent_team.nodes import review_loop
from agent_team.nodes.build_verify_subgraph import (
make_build_node,
make_verify_node,
route_after_verify,
)
from agent_team.nodes.verifier import VerifierConfig
review_loop.set_review_invoker(lambda prompt, **kw: review_text)
def fake_builder(*, plan, config):
return _p3_diff()
build_node = make_build_node(diff_builder=fake_builder)
verify_node = make_verify_node(
VerifierConfig(expected_run_id="r1", allowed_scope=["src"]),
ci_result_fetcher=ci_result_fetcher,
)
return build_graph(
checkpointer=_Saver(),
live_plan_node=_p3_plan_stub,
review_node=review_loop.bind_review_node(),
route_review=review_loop.route_after_review,
build_verify=(build_node, verify_node, route_after_verify),
)
def test_build_graph_build_verify_requires_review_node() -> None:
"""build_verify without review_node is a wiring error (no 'build' route)."""
from agent_team.nodes.build_verify_subgraph import (
make_build_node,
make_verify_node,
route_after_verify,
)
from agent_team.nodes.verifier import VerifierConfig
tuple_ = (
make_build_node(diff_builder=None),
make_verify_node(VerifierConfig(expected_run_id="")),
route_after_verify,
)
with pytest.raises(ValueError, match="build_verify"):
build_graph(build_verify=tuple_)
def test_p3_graph_authenticated_pass_routes_to_done(restore_review_invoker) -> None:
"""review(APPROVE) -> build -> verify(PASS via fake CI) -> DONE (PR terminus)."""
def pass_fetcher(state):
from agent_team.state_store import compute_content_hash
diff_hash = compute_content_hash(_p3_diff().encode("utf-8"))
return {"run_id": "r1", "conclusion": "success", "diff_hash": diff_hash}
graph = _p3_graph("VERDICT: APPROVE\nlooks solid", ci_result_fetcher=pass_fetcher)
thread_id, _ = start_task(graph, transport="slack")
final = resume_task(graph, thread_id=thread_id, answer="scope is X")
# The authenticated CI pass cleared the gate -> DONE terminus.
assert final["current_phase"] == Phase.DONE.value
assert final["status"] == TaskStatus.DONE.value
def test_p3_graph_inert_default_parks_at_verify(restore_review_invoker) -> None:
"""review(APPROVE) -> build -> verify(no CI result) -> BLOCK -> PARKED.
With the INERT default (no authenticated CI result) the gate can never
fabricate a pass, so an approved plan still parks at VERIFY. This is the
production-safe behavior the opt-in subgraph ships with.
"""
graph = _p3_graph(
"VERDICT: APPROVE\nlooks solid", ci_result_fetcher=lambda state: None
)
thread_id, _ = start_task(graph, transport="slack")
final = resume_task(graph, thread_id=thread_id, answer="scope is X")
assert final["current_phase"] == Phase.PARKED.value
assert final["status"] == TaskStatus.PARKED.value
def test_p3_graph_route_constants_mirror_subgraph_by_value() -> None:
"""graph.py's P3 route ids match the subgraph module by value (no cycle)."""
from agent_team.nodes import build_verify_subgraph as bvs
assert graph_mod.APPROVED_ROUTE == bvs.APPROVED_ROUTE
assert graph_mod.BUILD_ROUTE == bvs.BUILD_ROUTE
assert graph_mod.PARKED_ROUTE == bvs.PARKED_ROUTE
# --- Module import hygiene. ------------------------------------------------- # --- Module import hygiene. -------------------------------------------------

View file

@ -705,20 +705,154 @@ def test_start_runs_setup_and_start_task_and_prints_thread_id(
assert callable(coord.review_wiring) assert callable(coord.review_wiring)
def test_build_transport_live_github_raises_system_exit(cli: ModuleType) -> None: def test_build_transport_live_github_builds_github_transport(
"""A non-slack live transport is not wired and raises a clear SystemExit.""" cli: ModuleType, monkeypatch: pytest.MonkeyPatch
) -> None:
"""``--transport github`` builds the live GitHub transport from the env thread.
Token resolution is deferred to call time, so with GITHUB_TOKEN set and the
issue thread configured the builder constructs a GitHubTransport bound to
that ``owner/repo#issue`` (no network — only the issue thread is wired).
"""
monkeypatch.setenv("GITHUB_TOKEN", "ghp_test")
monkeypatch.setenv("GITHUB_OWNER", "Sea-Haven-Industries")
monkeypatch.setenv("GITHUB_REPO", "orchestrator")
monkeypatch.setenv("GITHUB_ISSUE_NUMBER", "42")
from agent_team.transport.github_adapter import GitHubTransport
args = argparse.Namespace(dry_run=False, transport="github") args = argparse.Namespace(dry_run=False, transport="github")
with pytest.raises(SystemExit, match="is not wired for the run-team CLI"): transport = cli._build_transport(args)
assert isinstance(transport, GitHubTransport)
assert transport.owner == "Sea-Haven-Industries"
assert transport.repo == "orchestrator"
assert transport.issue_number == 42
def test_build_transport_live_github_missing_thread_raises_system_exit(
cli: ModuleType, monkeypatch: pytest.MonkeyPatch
) -> None:
"""A github transport with no issue thread configured fails loudly."""
monkeypatch.delenv("GITHUB_OWNER", raising=False)
monkeypatch.delenv("GITHUB_REPO", raising=False)
monkeypatch.delenv("GITHUB_ISSUE_NUMBER", raising=False)
args = argparse.Namespace(dry_run=False, transport="github")
with pytest.raises(SystemExit, match="GITHUB_OWNER"):
cli._build_transport(args) cli._build_transport(args)
def test_build_transport_live_claude_code_raises_system_exit(cli: ModuleType) -> None: def test_build_transport_live_github_non_integer_issue_raises_system_exit(
"""claude_code is likewise un-wired for the P1 CLI surface.""" cli: ModuleType, monkeypatch: pytest.MonkeyPatch
) -> None:
"""A non-integer GITHUB_ISSUE_NUMBER fails loudly rather than half-building."""
monkeypatch.setenv("GITHUB_OWNER", "o")
monkeypatch.setenv("GITHUB_REPO", "r")
monkeypatch.setenv("GITHUB_ISSUE_NUMBER", "not-a-number")
args = argparse.Namespace(dry_run=False, transport="github")
with pytest.raises(SystemExit, match="must be an integer"):
cli._build_transport(args)
def test_build_transport_live_claude_code_builds_file_drop_transport(
cli: ModuleType, monkeypatch: pytest.MonkeyPatch, tmp_path: Path
) -> None:
"""``--transport claude_code`` builds the live file-drop transport from env."""
monkeypatch.setenv("CLAUDE_CODE_DROP_DIR", str(tmp_path / "drops"))
from agent_team.transport.claude_code_adapter import ClaudeCodeAdapter
args = argparse.Namespace(dry_run=False, transport="claude_code") args = argparse.Namespace(dry_run=False, transport="claude_code")
with pytest.raises(SystemExit, match="is not wired for the run-team CLI"): transport = cli._build_transport(args)
assert isinstance(transport, ClaudeCodeAdapter)
def test_build_transport_live_claude_code_missing_drop_dir_raises_system_exit(
cli: ModuleType, monkeypatch: pytest.MonkeyPatch
) -> None:
"""claude_code with no CLAUDE_CODE_DROP_DIR fails loudly."""
monkeypatch.delenv("CLAUDE_CODE_DROP_DIR", raising=False)
args = argparse.Namespace(dry_run=False, transport="claude_code")
with pytest.raises(SystemExit, match="CLAUDE_CODE_DROP_DIR"):
cli._build_transport(args) cli._build_transport(args)
def test_intake_github_polls_and_starts_tasks_dry_run(
cli: ModuleType,
db_path: Path,
audit_log: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""``intake-github --dry-run`` builds a coordinator, polls, and ingests issues.
The lazily-imported ``Coordinator`` is replaced with the fake (so no graph /
SDK is built), and ``build_default_issue_client`` is patched to return an
in-memory fake lister (so no GitHub network). The fake coordinator records
the ``start_task`` calls the intake leaf makes — one per labeled issue.
"""
_FakeCoordinator.instances.clear()
monkeypatch.setattr(
"agent_team.coordinator.Coordinator", _FakeCoordinator, raising=True
)
class _FakeIssueClient:
def list_open_issues(self, *, label: str) -> list[dict[str, Any]]:
assert label == "agent-team"
return [
{"id": 1001, "title": "Do thing A", "body": "details A"},
{"id": 1002, "title": "Do thing B", "body": ""},
]
def _fake_build_client(*, owner: str, repo: str, **_kw: Any) -> Any:
assert owner == "Sea-Haven-Industries"
assert repo == "orchestrator"
return _FakeIssueClient()
monkeypatch.setattr(
"agent_team.transport.github_intake.build_default_issue_client",
_fake_build_client,
raising=True,
)
code, out = _run(
cli,
db_path,
audit_log,
"intake-github",
"--dry-run",
"--owner",
"Sea-Haven-Industries",
"--repo",
"orchestrator",
"--label",
"agent-team",
)
assert code == 0
assert len(_FakeCoordinator.instances) == 1
coord = _FakeCoordinator.instances[0]
assert coord.setup_called is True
# Both labeled issues were ingested; the leaf passes title+body and the
# GitHub transport name. (The fake records only the LAST call's kwargs.)
assert coord.start_kwargs == {
"task_text": "Do thing B",
"transport_name": "github",
}
# The ingested issue ids are printed (one per line).
assert out.split() == ["1001", "1002"]
def test_intake_github_label_required(cli: ModuleType) -> None:
"""``intake-github`` requires --owner/--repo/--label (argparse usage error)."""
parser = cli.build_parser()
with pytest.raises(SystemExit):
parser.parse_args(["intake-github", "--owner", "o", "--repo", "r"])
def test_build_transport_dry_run_returns_dry_run_transport(cli: ModuleType) -> None: def test_build_transport_dry_run_returns_dry_run_transport(cli: ModuleType) -> None:
"""``dry_run=True`` yields a _DryRunTransport whose post returns a synthetic ref.""" """``dry_run=True`` yields a _DryRunTransport whose post returns a synthetic ref."""
args = argparse.Namespace(dry_run=True, transport="slack") args = argparse.Namespace(dry_run=True, transport="slack")