From 082e45bf881a83330de769c2a93beddca768d798 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 20:40:51 -0400 Subject: [PATCH] test(agent-team): end-to-end plan-review gate composition (Phase C) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Hermetic e2e tests driving the WHOLE stack composed together — real build_graph(plan_gate=True) + real Coordinator + real SQLite ledger + real review_loop router, with only the LLM nodes stubbed — through the daemon API (start_task/submit_answer/tick), never nodes directly. Flows: approve settles at BUILD; request-changes via RAW PROSE through the real SlackListener -> _resolve_payload -> map_plan_decision proves the prose maps to request_changes (NOT FAILED) and the notes reach the planner; abandon -> FAILED; repeated request_changes terminates at MAX_PLAN_GATE_VISITS -> PARKED; and a legacy (no-kind) ledger migrates in place then routes clarify vs plan_decision correctly. 1459 passed. --- agent-team/tests/test_plan_gate_e2e.py | 544 +++++++++++++++++++++++++ 1 file changed, 544 insertions(+) create mode 100644 agent-team/tests/test_plan_gate_e2e.py diff --git a/agent-team/tests/test_plan_gate_e2e.py b/agent-team/tests/test_plan_gate_e2e.py new file mode 100644 index 0000000..28a526e --- /dev/null +++ b/agent-team/tests/test_plan_gate_e2e.py @@ -0,0 +1,544 @@ +"""Phase C end-to-end integration tests for the plan-review decision gate. + +These drive a task through the WHOLE Plane-2 stack composed together — the real +:func:`agent_team.graph.build_graph` (with ``plan_gate=True``), the real +:class:`agent_team.coordinator.Coordinator`, a real SQLite ledger (tmp), an +in-memory LangGraph checkpointer, and (for the kind-aware path) the real +:class:`agent_team.transport.slack_listener.SlackListener` + +:func:`agent_team.transport.slack_adapter.map_plan_decision`. ONLY the three LLM +nodes are stubbed so the run is deterministic and hermetic (no live model, no +Slack network, no GitHub): + +* the **clarifier** is the deterministic single-turn ``graph.clarify_node`` stub + (it really suspends on ``interrupt()`` — the human gate is real); +* the **planner** is a stub that emits a plan and advances to REVIEW (mirroring + ``planner.plan_node``'s contract), reused from the committed graph tests; +* the **reviewer** is the REAL ``review_loop`` node + router bound to a stubbed + review invoker (a callable returning canned ``VERDICT:`` text), so the + plan<->review loop, the round cap, and the cap→plan-gate dead-end are all real. + +Everything else — INTAKE, the human-gate suspend/resume mechanic, the +``plan_gate_node`` interrupt + decision parsing + visit ceiling, the +coordinator's ledger row opening / single-open-gate invariant / presentation +threading / kind routing, the listener's Slack→decision composition, and the +schema migration — is the real production code path. + +The drive surface is the daemon's public API (``start_task`` / ``submit_answer`` +/ ``tick`` which drains + posts follow-ups), NEVER the nodes directly. +""" + +from __future__ import annotations + +import queue +from pathlib import Path +from typing import Any + +import pytest + +try: # InMemorySaver is the modern name; fall back on older langgraph. + from langgraph.checkpoint.memory import InMemorySaver as _Saver +except ImportError: # pragma: no cover - environment-dependent + from langgraph.checkpoint.memory import MemorySaver as _Saver + +from agent_team import graph as graph_mod +from agent_team.coordinator import Coordinator +from agent_team.db.schema import connect, init_db +from agent_team.nodes import review_loop +from agent_team.task_model import Phase, PipelineState, TaskStatus +from agent_team.transport.base import QuestionSet, Transport +from agent_team.transport.slack_adapter import SlackTransport +from agent_team.transport.slack_listener import SlackListener + +# --------------------------------------------------------------------------- # +# Test doubles + harness (composed from the committed coordinator/graph tests). +# --------------------------------------------------------------------------- # + +# The single authorized owner id (AUTHZ-01) used on the listener path. +OWNER_ID = "U_OWNER" + + +class FakeTransport(Transport): + """Record-only transport: no Slack, no network (the §3.3.1 injection seam). + + Mirrors the FakeTransport in ``tests/test_coordinator.py``: ``post_question`` + records the posted question-set + the thread_ts it was threaded under and + returns a deterministic ``channel_ref`` embedding the ``question_id``; + ``parse_answer`` reads a plain ``{"question_id","answer","via"}`` dict so a + decision can be submitted without a Slack payload. + """ + + def __init__(self) -> None: + self.posted: list[QuestionSet] = [] + self.thread_tss: list[str | None] = [] + + def post_question( + self, + *, + thread_id: str, + question_id: str, + turn: int, + question_set: QuestionSet, + deadline: str, + thread_ts: str | None = None, + ) -> str: + self.posted.append(question_set) + self.thread_tss.append(thread_ts) + return f"fake:{question_id}" + + def parse_answer(self, raw: Any) -> tuple[str, Any, str]: + return raw["question_id"], raw["answer"], raw.get("via", "fake") + + +def _plan_stub(state: PipelineState) -> PipelineState: + """Planner stub: emit a plan and advance to REVIEW (no model call). + + Mirrors ``planner.plan_node``'s contract (sets ``plan`` + phase REVIEW), + reused verbatim from the committed P2 graph tests. The revision index tracks + prior review rounds so a re-plan after request_changes is observable, and the + folded-in human notes (review feedback) can be asserted by a wrapping stub. + """ + revisions = len(state.get("review_verdicts") or []) + return PipelineState( + plan={ + "summary": "Add a hermetic plan-gate smoke test.", + "phases": [{"name": "Audit infra"}, {"name": "Write test"}], + "revision": revisions, + }, + current_phase=Phase.REVIEW.value, + status=TaskStatus.ACTIVE.value, + ) + + +def _e2e_coordinator( + db_path: Path, + *, + review_text: "str | Any", + transport: Transport | None = None, + resume_queue: "queue.Queue[Any] | None" = None, + plan_node: Any = None, + max_review_rounds: int = 1, + notify: Any = None, +) -> Coordinator: + """Compose the REAL graph + coordinator with only the LLM nodes stubbed. + + * clarifier = the real ``graph.clarify_node`` suspend stub; + * planner = ``_plan_stub`` (or an injected wrapping stub); + * reviewer = the REAL ``review_loop`` node + router, bound to a stubbed + review invoker (``review_text`` may be a string or a callable that + LangGraph-style ``invoker(prompt, **kw)`` consumes); + * ``plan_gate=True`` so the review-cap dead-end suspends on the resumable + ``plan_decision`` interrupt instead of terminally parking. + + A fresh in-memory saver + the tmp SQLite ledger make the durable composition + real. ``max_review_rounds=1`` reaches the gate after a single REQUEST CHANGES + so the loop is short and deterministic. + """ + if callable(review_text): + review_loop.set_review_invoker(review_text) + else: + review_loop.set_review_invoker(lambda prompt, **kw: review_text) + saver = _Saver() + return Coordinator( + db_path=db_path, + transport=transport or FakeTransport(), + build_clarify_node=lambda: graph_mod.clarify_node, + build_plan_node=lambda: plan_node or _plan_stub, + review_wiring=lambda: ( + review_loop.bind_review_node({"max_review_rounds": max_review_rounds}), + review_loop.route_after_review, + ), + build_checkpointer=lambda _path: saver, + resume_queue=resume_queue, + notify=notify, + plan_gate=True, + ) + + +@pytest.fixture() +def restore_review_invoker(): + """Save/restore the review-loop module-global invoker around each test.""" + saved = review_loop._review_invoker + yield + review_loop._review_invoker = saved + + +@pytest.fixture() +def db_path(tmp_path: Path) -> Path: + path = tmp_path / "state" / "agent_team.sqlite" + init_db(path) + return path + + +# --------------------------------------------------------------------------- # +# Ledger helpers. +# --------------------------------------------------------------------------- # + + +def _open_rows(db_path: Path) -> list[dict[str, Any]]: + conn = connect(db_path) + try: + rows = conn.execute( + "SELECT * FROM pending_questions WHERE status='open'" + ).fetchall() + finally: + conn.close() + return [dict(r) for r in rows] + + +def _rows_by_kind(db_path: Path, kind: str) -> list[dict[str, Any]]: + conn = connect(db_path) + try: + rows = conn.execute( + "SELECT * FROM pending_questions WHERE kind=?", (kind,) + ).fetchall() + finally: + conn.close() + return [dict(r) for r in rows] + + +def _open_clarifier_qid(db_path: Path) -> str: + rows = _open_rows(db_path) + assert len(rows) == 1 + return str(rows[0]["question_id"]) + + +def _drive_to_plan_gate( + coord: Coordinator, *, root_ts: str = "ROOT.1" +) -> tuple[str, dict[str, Any]]: + """Start a task, answer the clarifier, and tick to the plan gate. + + Returns ``(thread_id, gate_row)`` where ``gate_row`` is the opened + ``kind='plan_decision'`` ledger row. Exercises the real surface end to end: + ``start_task`` -> clarifier suspend -> ``submit_answer`` -> ``tick`` (drain + runs intake->clarify->plan->review loop to the cap, suspends on the gate, and + ``_post_resume_followups`` opens + presents the decision row). + """ + thread_id = coord.start_task( + task_text="add a smoke test", + transport_name="slack", + slack_thread_ts=root_ts, + ) + clar_qid = _open_clarifier_qid(db_path_of(coord)) + coord.submit_answer({"question_id": clar_qid, "answer": "scope is X", "via": "v"}) + coord.tick() + + gate_rows = _rows_by_kind(db_path_of(coord), graph_mod.PLAN_DECISION_KIND) + assert len(gate_rows) == 1, "exactly one plan_decision gate row must be opened" + return thread_id, gate_rows[0] + + +def db_path_of(coord: Coordinator) -> Path: + return coord._db_path + + +# --------------------------------------------------------------------------- # +# 1. Approve path. +# --------------------------------------------------------------------------- # + + +def test_e2e_approve_path_settles_approved( + db_path: Path, restore_review_invoker +) -> None: + """start -> clarify -> plan -> review(ESCALATE) -> gate opened+presented -> + APPROVE decision -> task settles in the approved terminal state, no open row. + + The reviewer never approves, so the plan<->review loop hits the cap and + suspends on the real ``plan_decision`` gate. The coordinator opens the durable + row and posts the presentation THREADED under the task root. An ``approve`` + decision (through ``submit_answer``) drives the real ``plan_gate_node`` to the + same approved terminus an auto-approved plan reaches (phase BUILD, ACTIVE). + """ + transport = FakeTransport() + coord = _e2e_coordinator( + db_path, review_text="VERDICT: REQUEST CHANGES\nnot ready", transport=transport + ) + coord.setup() + + thread_id, gate_row = _drive_to_plan_gate(coord) + + # The opened gate row is durable, open, kind='plan_decision', threaded under + # the task root (its channel_ref is the root ts so a reply maps back). + assert gate_row["kind"] == graph_mod.PLAN_DECISION_KIND + assert gate_row["status"] == "open" + assert gate_row["channel_ref"] == "ROOT.1" + # A presentation was posted threaded under the root, carrying the plan. + assert transport.thread_tss[-1] == "ROOT.1" + gate_qset = transport.posted[-1] + assert gate_qset.context.get("kind") == graph_mod.PLAN_DECISION_KIND + assert "Summary:" in gate_qset.context.get("presentation", "") + + # Approve the plan through the real submit_answer -> resume -> graph path. + coord.submit_answer( + { + "question_id": gate_row["question_id"], + "answer": {"decision": "approve", "notes": ""}, + "via": "v", + } + ) + coord.tick() + + state = graph_mod.get_pipeline_state(coord.graph, thread_id=thread_id) + assert state["current_phase"] == Phase.BUILD.value + assert state["status"] == TaskStatus.ACTIVE.value + # No gate stays open and no pending interrupt remains. + assert _open_rows(db_path) == [] + assert graph_mod.pending_question(coord.graph, thread_id=thread_id) is None + + +# --------------------------------------------------------------------------- # +# 2. Request-changes path — RAW PROSE through the kind-aware listener. +# --------------------------------------------------------------------------- # + + +def test_e2e_request_changes_prose_via_listener_loops_to_planner( + db_path: Path, restore_review_invoker +) -> None: + """A FREE-TEXT prose reply, routed through the REAL SlackListener + + map_plan_decision, maps to request_changes (NOT abandon/FAILED) and loops the + task back to the planner with the human notes folded into review feedback. + + This is the load-bearing composition assertion: the prose + ("please use pytest fixtures") goes through ``handle_event`` -> + ``_resolve_payload`` -> ``map_plan_decision`` (the kind-aware Slack→decision + map), is enqueued onto the coordinator's resume queue, and the next ``tick`` + re-plans. Without the kind-aware map, the graph's decision parser would map + the prose verb to ``abandon`` -> terminal FAILED. + """ + # First review round REQUEST CHANGES (to reach the gate); after the human's + # prose request_changes loops back, the next review APPROVES so the task + # settles rather than re-gating (keeps the assertion crisp). + texts = iter(["VERDICT: REQUEST CHANGES\nnot ready", "VERDICT: APPROVE\nnow good"]) + last = {"v": "VERDICT: APPROVE\nnow good"} + + def review_invoker(prompt, **kw): + try: + last["v"] = next(texts) + except StopIteration: + pass + return last["v"] + + # A wrapping planner stub that captures the review feedback the planner sees + # on re-plan, so we can prove the human notes reached it. + from agent_team.nodes.planner import _format_review_feedback + + captured: dict[str, Any] = {} + + def capturing_plan(state: PipelineState) -> PipelineState: + captured["feedback"] = _format_review_feedback( + list(state.get("review_verdicts") or []) + ) + return _plan_stub(state) + + transport = FakeTransport() + # The coordinator and the listener SHARE one resume queue: the listener + # enqueues onto it, the coordinator drains it (the production handoff seam). + shared_q: "queue.Queue[Any]" = queue.Queue() + coord = _e2e_coordinator( + db_path, + review_text=review_invoker, + transport=transport, + resume_queue=shared_q, + plan_node=capturing_plan, + ) + coord.setup() + + thread_id, gate_row = _drive_to_plan_gate(coord) + # The gate post threaded under the root; its channel_ref is the root ts, so a + # thread reply with thread_ts == root ts maps back to THIS gate row. + root_ts = gate_row["channel_ref"] + assert root_ts == "ROOT.1" + + # Build the REAL listener over the SAME ledger DB, enqueueing onto the SAME + # queue the coordinator drains. The listener uses a real SlackTransport + # (parse_answer + channel only; no network). + listener = SlackListener( + SlackTransport(channel="C123"), + db_path, + shared_q.put, + owner_ids={OWNER_ID}, + ) + + # A REAL slack_bolt Events API thread-reply envelope carrying arbitrary + # change-request PROSE (no callback_id / question_id / metadata): the question + # is resolved by thread_ts == the gate row's channel_ref, and because that row + # is kind='plan_decision' the prose is mapped via map_plan_decision. + prose_reply = { + "type": "event_callback", + "event": { + "type": "message", + "text": "please use pytest fixtures", + "thread_ts": root_ts, + "ts": "1700000001.000200", + "channel": "C123", + "user": OWNER_ID, + }, + } + outcome = listener.handle_event(prose_reply) + assert outcome is not None and outcome.accepted is True + + # Drain the listener-enqueued resume through the coordinator: the graph loops + # plan -> review (now APPROVE) and settles. It must NOT have FAILED/abandoned. + coord.tick() + + state = graph_mod.get_pipeline_state(coord.graph, thread_id=thread_id) + assert state["status"] != TaskStatus.FAILED.value, ( + "raw prose must map to request_changes, never an accidental abandon/FAILED" + ) + # The human's exact prose reached the planner's review-feedback formatter on + # the re-plan (proving request_changes folded the notes in). + assert "please use pytest fixtures" in str(captured.get("feedback", "")) + # And the loop re-entered plan -> review and settled past the gate. + assert state["current_phase"] == Phase.BUILD.value + assert _open_rows(db_path) == [] + + +# --------------------------------------------------------------------------- # +# 3. Abandon path. +# --------------------------------------------------------------------------- # + + +def test_e2e_abandon_path_fails(db_path: Path, restore_review_invoker) -> None: + """An ``abandon`` decision at the gate drives the task to terminal FAILED.""" + coord = _e2e_coordinator(db_path, review_text="VERDICT: REQUEST CHANGES\nnope") + coord.setup() + + thread_id, gate_row = _drive_to_plan_gate(coord) + coord.submit_answer( + { + "question_id": gate_row["question_id"], + "answer": {"decision": "abandon", "notes": ""}, + "via": "v", + } + ) + coord.tick() + + state = graph_mod.get_pipeline_state(coord.graph, thread_id=thread_id) + assert state["status"] == TaskStatus.FAILED.value + assert _open_rows(db_path) == [] + + +# --------------------------------------------------------------------------- # +# 4. Ceiling termination — repeated request_changes terminates PARKED. +# --------------------------------------------------------------------------- # + + +def test_e2e_repeated_request_changes_hits_ceiling_parked( + db_path: Path, restore_review_invoker +) -> None: + """Repeated request_changes eventually hits ``MAX_PLAN_GATE_VISITS`` and goes + terminal PARKED — the gate loop does NOT spin forever. + + The reviewer NEVER approves, so every re-plan re-hits the review cap and + re-suspends on the gate; each request_changes consumes one gate visit. The + test loops well above the ceiling and asserts it terminates on its own (no + open gate, status PARKED, ceiling reason) at exactly the visit cap. + """ + coord = _e2e_coordinator( + db_path, review_text="VERDICT: REQUEST CHANGES\nstill not ready" + ) + coord.setup() + + thread_id, gate_row = _drive_to_plan_gate(coord) + + # Each iteration: answer the open gate request_changes, then tick (re-plan -> + # review cap -> re-suspend OR terminate). Bound the loop well above the + # ceiling to prove it self-terminates, not via our cap. + for _ in range(graph_mod.MAX_PLAN_GATE_VISITS + 5): + open_rows = _rows_by_kind(db_path, graph_mod.PLAN_DECISION_KIND) + open_now = [r for r in open_rows if r["status"] == "open"] + if not open_now: + break + coord.submit_answer( + { + "question_id": open_now[0]["question_id"], + "answer": {"decision": "request_changes", "notes": "again"}, + "via": "v", + } + ) + coord.tick() + + state = graph_mod.get_pipeline_state(coord.graph, thread_id=thread_id) + assert graph_mod.pending_question(coord.graph, thread_id=thread_id) is None + assert state["status"] == TaskStatus.PARKED.value + assert "ceiling" in (state.get("failure_reason") or "") + assert state.get("plan_gate_visits") == graph_mod.MAX_PLAN_GATE_VISITS + assert _open_rows(db_path) == [] + + +# --------------------------------------------------------------------------- # +# 5. Migration mixed-state composition — legacy clarifier row + new gate row. +# --------------------------------------------------------------------------- # + + +# Legacy (pre-kind) pending_questions DDL, mirrored from tests/test_schema.py, to +# construct a DB whose table predates the additive ``kind`` migration. +_LEGACY_PENDING_QUESTIONS_DDL = """ +CREATE TABLE IF NOT EXISTS pending_questions ( + question_id TEXT PRIMARY KEY, + thread_id TEXT NOT NULL, + turn INTEGER NOT NULL, + status TEXT NOT NULL + CHECK (status IN ('open', 'answered', 'expired', 'superseded')), + transport TEXT NOT NULL, + channel_ref TEXT, + posted_at TEXT, + deadline_at TEXT, + answer_json TEXT, + answered_at TEXT, + answered_via TEXT +) +""".strip() + + +def test_e2e_migration_mixed_state_legacy_clarify_plus_new_gate( + tmp_path: Path, restore_review_invoker +) -> None: + """A coordinator pointed at an OLD-schema ledger (no ``kind`` column) carrying + a legacy clarifier row migrates it, then drives a NEW task to the plan gate. + + Proves the migration + kind routing COMPOSE on a real upgraded ledger: + + * the pre-existing legacy row reads back ``kind='clarify'`` (the migration + default), and + * the new gate row the coordinator opens is ``kind='plan_decision'`` (the + kind discriminator routes correctly on the upgraded DB). + """ + db_path = tmp_path / "legacy" / "agent_team.sqlite" + db_path.parent.mkdir(parents=True, exist_ok=True) + + # 1. Build the OLD table by hand (no kind column) and seed a legacy clarifier + # row, then close — simulating a pre-upgrade ledger on disk. + conn = connect(db_path) + try: + conn.execute(_LEGACY_PENDING_QUESTIONS_DDL) + conn.execute( + "INSERT INTO pending_questions " + "(question_id, thread_id, turn, status, transport) " + "VALUES ('legacy-clarify', 'legacy-thread', 0, 'answered', 'slack')" + ) + assert "kind" not in { + r["name"] for r in conn.execute("PRAGMA table_info(pending_questions)") + } + finally: + conn.close() + + # 2. Coordinator.setup() runs init_db (the migration) on the existing DB, then + # we drive a NEW task to the plan gate over the upgraded ledger. + coord = _e2e_coordinator(db_path, review_text="VERDICT: REQUEST CHANGES\nnope") + coord.setup() # migrates the legacy table in place (adds kind column) + + _thread_id, gate_row = _drive_to_plan_gate(coord) + + # The legacy row survived the migration and reads the 'clarify' default. + conn = connect(db_path) + try: + legacy_kind = conn.execute( + "SELECT kind FROM pending_questions WHERE question_id='legacy-clarify'" + ).fetchone()["kind"] + finally: + conn.close() + assert legacy_kind == "clarify" + + # And the new gate row the coordinator opened is kind='plan_decision' — the + # migration + kind routing compose on a real upgraded ledger. + assert gate_row["kind"] == graph_mod.PLAN_DECISION_KIND