agent-team Plane-2: P4 live transports + intake, P3-inert build/verify subgraph #13
13 changed files with 2892 additions and 21 deletions
|
|
@ -66,6 +66,7 @@ if TYPE_CHECKING: # pragma: no cover - typing only
|
|||
|
||||
__all__ = [
|
||||
"Coordinator",
|
||||
"build_verify_wiring",
|
||||
"default_clarify_node_factory",
|
||||
]
|
||||
|
||||
|
|
@ -96,6 +97,19 @@ ReviewWiring = Callable[
|
|||
[],
|
||||
"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
|
||||
# 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
|
||||
|
||||
|
||||
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:
|
||||
"""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_plan_node: PlanNodeFactory | None = None,
|
||||
review_wiring: ReviewWiring | None = None,
|
||||
build_verify_wiring: BuildVerifyWiring | None = None,
|
||||
build_checkpointer: CheckpointerFactory | None = None,
|
||||
resume_queue: "queue.Queue[Any] | None" = None,
|
||||
deadline_window: timedelta | None = None,
|
||||
|
|
@ -224,6 +290,10 @@ class Coordinator:
|
|||
# path injects default_plan_node_factory + default_review_wiring.
|
||||
self._build_plan_node = build_plan_node
|
||||
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 = (
|
||||
build_checkpointer or graph_mod.build_sqlite_checkpointer
|
||||
)
|
||||
|
|
@ -294,12 +364,20 @@ class Coordinator:
|
|||
if self._review_wiring is not None:
|
||||
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(
|
||||
checkpointer,
|
||||
live_clarify_node=clarify_node,
|
||||
live_plan_node=plan_node,
|
||||
review_node=review_node,
|
||||
route_review=route_review,
|
||||
build_verify=build_verify,
|
||||
)
|
||||
|
||||
# The ResumeWorker is satisfied directly by the compiled LangGraph app
|
||||
|
|
@ -354,7 +432,14 @@ class Coordinator:
|
|||
if self._graph is None:
|
||||
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)
|
||||
|
||||
question = graph_mod.pending_question(self._graph, thread_id=thread_id)
|
||||
|
|
|
|||
|
|
@ -62,6 +62,8 @@ if TYPE_CHECKING: # pragma: no cover - typing only
|
|||
from langgraph.graph.state import CompiledStateGraph
|
||||
|
||||
__all__ = [
|
||||
"APPROVED_ROUTE",
|
||||
"BUILD_NODE",
|
||||
"BUILD_ROUTE",
|
||||
"CLARIFY",
|
||||
"DEFAULT_CLARIFY_DEADLINE",
|
||||
|
|
@ -70,6 +72,7 @@ __all__ = [
|
|||
"PARKED_ROUTE",
|
||||
"PLAN",
|
||||
"REVIEW",
|
||||
"VERIFY_NODE",
|
||||
"build_graph",
|
||||
"build_sqlite_checkpointer",
|
||||
"clarify_node",
|
||||
|
|
@ -99,6 +102,25 @@ REVIEW = "review"
|
|||
BUILD_ROUTE = "build"
|
||||
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
|
||||
# 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").
|
||||
|
|
@ -266,6 +288,12 @@ def build_graph(
|
|||
live_plan_node: Callable[[PipelineState], PipelineState] | None = None,
|
||||
review_node: Callable[[PipelineState], PipelineState] | 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:
|
||||
"""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
|
||||
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
|
||||
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)."
|
||||
)
|
||||
|
||||
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.add_node(INTAKE, intake_node)
|
||||
builder.add_node(CLARIFY, clarify)
|
||||
|
|
@ -334,14 +390,38 @@ def build_graph(
|
|||
# P1: the plan stage is the terminus.
|
||||
builder.add_edge(PLAN, END)
|
||||
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_edge(PLAN, REVIEW)
|
||||
builder.add_conditional_edges(
|
||||
REVIEW,
|
||||
route_review,
|
||||
{BUILD_ROUTE: END, PLAN: PLAN, PARKED_ROUTE: END},
|
||||
)
|
||||
|
||||
if build_verify is None:
|
||||
# P2: the review's "build" route is the approved-plan terminus.
|
||||
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:
|
||||
return builder.compile()
|
||||
|
|
|
|||
286
agent-team/agent_team/nodes/build_verify_subgraph.py
Normal file
286
agent-team/agent_team/nodes/build_verify_subgraph.py
Normal 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}"
|
||||
)
|
||||
265
agent-team/agent_team/transport/claude_code_live.py
Normal file
265
agent-team/agent_team/transport/claude_code_live.py
Normal 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
|
||||
293
agent-team/agent_team/transport/github_intake.py
Normal file
293
agent-team/agent_team/transport/github_intake.py
Normal 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()
|
||||
229
agent-team/agent_team/transport/github_live.py
Normal file
229
agent-team/agent_team/transport/github_live.py
Normal 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)
|
||||
|
|
@ -38,6 +38,9 @@ Subcommands (P1 surface):
|
|||
``--confirm``; audit-logged). Maps to :func:`answer_question`.
|
||||
* ``supersede`` — mark a stale question ``superseded`` (DESTRUCTIVE: requires
|
||||
``--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
|
||||
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)
|
||||
|
||||
# Transport choices the start/serve commands accept (§3.3.1 D10). Only ``slack``
|
||||
# has a live adapter wired for the P1 CLI; the others are accepted for forward
|
||||
# compatibility and gated in _build_transport.
|
||||
# Transport choices the start/serve/intake commands accept (§3.3.1 D10). All
|
||||
# three now have a live human-gate adapter wired in _build_transport: ``slack``
|
||||
# (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")
|
||||
|
||||
__all__ = [
|
||||
|
|
@ -511,10 +517,23 @@ def _build_transport(args: argparse.Namespace) -> Any:
|
|||
"""Build the transport for a coordinator command (lazy; token-tolerant).
|
||||
|
||||
``--dry-run`` (or any transport in dry-run) yields a non-posting transport so
|
||||
intake works without credentials. Otherwise the live Slack transport is
|
||||
constructed lazily from ``SLACK_BOT_TOKEN`` / ``SLACK_CHANNEL``; GitHub and
|
||||
Claude-Code live transports are not wired for the P1 CLI surface and raise a
|
||||
clear error rather than pretending to post.
|
||||
intake works without credentials. Otherwise the live transport is built
|
||||
lazily from the environment for the chosen ``--transport`` (so import,
|
||||
``--help``, and ledger commands never need a token):
|
||||
|
||||
* ``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):
|
||||
return _DryRunTransport()
|
||||
|
|
@ -523,12 +542,74 @@ def _build_transport(args: argparse.Namespace) -> Any:
|
|||
|
||||
channel = os.environ.get("SLACK_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(
|
||||
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):
|
||||
"""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
|
||||
|
||||
|
||||
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:
|
||||
"""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_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
|
||||
|
||||
|
||||
|
|
|
|||
375
agent-team/tests/test_build_verify_subgraph.py
Normal file
375
agent-team/tests/test_build_verify_subgraph.py
Normal 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
|
||||
280
agent-team/tests/test_claude_code_live.py
Normal file
280
agent-team/tests/test_claude_code_live.py
Normal 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
|
||||
233
agent-team/tests/test_github_intake.py
Normal file
233
agent-team/tests/test_github_intake.py
Normal 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
|
||||
333
agent-team/tests/test_github_live.py
Normal file
333
agent-team/tests/test_github_live.py
Normal 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"
|
||||
|
|
@ -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
|
||||
|
||||
|
||||
# --- 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. -------------------------------------------------
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -705,20 +705,154 @@ def test_start_runs_setup_and_start_task_and_prints_thread_id(
|
|||
assert callable(coord.review_wiring)
|
||||
|
||||
|
||||
def test_build_transport_live_github_raises_system_exit(cli: ModuleType) -> None:
|
||||
"""A non-slack live transport is not wired and raises a clear SystemExit."""
|
||||
def test_build_transport_live_github_builds_github_transport(
|
||||
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")
|
||||
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)
|
||||
|
||||
|
||||
def test_build_transport_live_claude_code_raises_system_exit(cli: ModuleType) -> None:
|
||||
"""claude_code is likewise un-wired for the P1 CLI surface."""
|
||||
def test_build_transport_live_github_non_integer_issue_raises_system_exit(
|
||||
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")
|
||||
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)
|
||||
|
||||
|
||||
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:
|
||||
"""``dry_run=True`` yields a _DryRunTransport whose post returns a synthetic ref."""
|
||||
args = argparse.Namespace(dry_run=True, transport="slack")
|
||||
|
|
|
|||
Reference in a new issue