diff --git a/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index 79d8ed5..bfe0760 100644 --- a/agent-team/agent_team/coordinator.py +++ b/agent-team/agent_team/coordinator.py @@ -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) diff --git a/agent-team/agent_team/graph.py b/agent-team/agent_team/graph.py index 8f6ee4a..685ee5c 100644 --- a/agent-team/agent_team/graph.py +++ b/agent-team/agent_team/graph.py @@ -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() diff --git a/agent-team/agent_team/nodes/build_verify_subgraph.py b/agent-team/agent_team/nodes/build_verify_subgraph.py new file mode 100644 index 0000000..fba59ec --- /dev/null +++ b/agent-team/agent_team/nodes/build_verify_subgraph.py @@ -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}" +) diff --git a/agent-team/agent_team/transport/claude_code_live.py b/agent-team/agent_team/transport/claude_code_live.py new file mode 100644 index 0000000..3211f90 --- /dev/null +++ b/agent-team/agent_team/transport/claude_code_live.py @@ -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 ```` 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 diff --git a/agent-team/agent_team/transport/github_intake.py b/agent-team/agent_team/transport/github_intake.py new file mode 100644 index 0000000..a3db56f --- /dev/null +++ b/agent-team/agent_team/transport/github_intake.py @@ -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=``, ``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 #`` 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=, + 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=