"""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