From ba4fe68fdb10778c992d61ac88df79c92a5be896 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 20:19:51 -0400 Subject: [PATCH] feat(agent-team): wire the plan-review gate into the coordinator (Phase B2b) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Connect the graph plan-gate (B2a) to the durable ledger + Slack presentation. - setup() passes build_graph(plan_gate=True) only on the wired review path (plan_gate flag ANDed with review_node present); P1/stub paths force it off. - _post_resume_followups detects a settled interrupt by payload kind == PLAN_DECISION_KIND (NOT status, which still reads 'parked' at the gate per B2a) and posts the decision gate: opens a pending_questions row with kind='plan_decision' (24h deadline, threaded, channel_ref = root ts) and presents _summarize_plan + findings + reply instructions, truncated to a ~2700-char Slack budget. A clarify/legacy interrupt keeps the existing path. - Single-open-gate invariant: the opener skips if any open row already exists for the thread (one row, one presentation). - Expiry: _park posts a plan-decision-specific recovery notice (re-assign / force-resume) for an expired gate row. - Resume path unchanged: the decision answer flows through submit_answer → ResumeWorker → plan_gate_node with no resume-worker special-casing. notify_question has no kind param in this tree, so the opener calls ledger.post_question(kind=...) + ledger.set_channel_ref directly; the clarifier path still uses notify_question unchanged. Tests: gate row+presentation+threading, approve/request_changes/abandon via submit_answer, single-open-gate skip, expiry notice. 1412 passed. --- agent-team/agent_team/coordinator.py | 354 ++++++++++++++++++++++++++- agent-team/tests/test_coordinator.py | 295 ++++++++++++++++++++++ 2 files changed, 645 insertions(+), 4 deletions(-) diff --git a/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index 9bea15b..c9424a7 100644 --- a/agent-team/agent_team/coordinator.py +++ b/agent-team/agent_team/coordinator.py @@ -67,6 +67,7 @@ from agent_team.transport.slack_adapter import SlackTransport if TYPE_CHECKING: # pragma: no cover - typing only from agent_team.task_model import PipelineState + from agent_team.transport.base import QuestionSet __all__ = [ "Coordinator", @@ -149,6 +150,26 @@ AlarmHook = Callable[[str], None] ListenerFactory = Callable[[], Any] +def _open_question_id_for_thread(conn: Any, thread_id: str) -> str | None: + """Return the ``question_id`` of an open ledger row for ``thread_id``, or None. + + The single-open-gate invariant (B-2) holds that a thread has at most one open + ``pending_questions`` row at a time. This is the read side of that guard: a + non-None result means an open gate already exists for the thread, so a second + gate (clarifier follow-up OR plan-decision) must NOT be opened. Returns the + first open row's id (there should never be more than one) or ``None``. + """ + if not thread_id: + return None + row = conn.execute( + "SELECT question_id FROM pending_questions " + "WHERE thread_id=? AND status='open' " + "ORDER BY posted_at ASC, question_id ASC LIMIT 1", + (thread_id,), + ).fetchone() + return None if row is None else str(row["question_id"]) + + def default_slack_listener_factory( *, transport: SlackTransport, @@ -582,6 +603,7 @@ class Coordinator: ci_poller: "Callable[[Any], Any] | None" = None, ci_timeout: timedelta | None = None, draft_pr_provider: "Callable[[], list[Any]] | None" = None, + plan_gate: bool = True, ) -> None: self._db_path = Path(db_path) self._transport = transport @@ -639,6 +661,16 @@ class Coordinator: # re-ALARMed / re-reminded every tick. self._draft_pr_provider = draft_pr_provider self._draft_pr_memory: Any = None + # Plan-review human decision gate (Phase B2b). When True AND a review loop + # is wired (``review_wiring`` is not None), ``setup`` builds the graph with + # ``plan_gate=True`` so the review-cap dead-end suspends on a resumable + # human decision instead of terminally parking. When the review loop is + # NOT wired (P1-only / stub paths, e.g. the unit suite's default + # coordinator) the gate cannot exist — the graph would raise — so setup + # forces it off regardless of this flag. Defaulting to True makes the live + # serve/P2 path gate-enabled out of the box; tests build both shapes by + # toggling this with/without ``review_wiring``. + self._plan_gate = plan_gate # Built by setup(). self._graph: Any = None @@ -744,6 +776,14 @@ class Coordinator: # transition_recorder=None (no instrumentation). transition_recorder = TransitionRecorder(self._db_path) + # Plan-review human decision gate (B2b): only enable it on the wired + # review path. ``build_graph`` raises if ``plan_gate`` is set without a + # ``review_node`` (the gate IS the review-cap dead-end's replacement), so + # gate off whenever the review loop is absent (P1-only / stub paths) even + # if the ctor flag asked for it. On the live serve/P2 path the review + # wiring is injected, so the gate turns on. + plan_gate = self._plan_gate and review_node is not None + self._graph = graph_mod.build_graph( checkpointer, transition_recorder=transition_recorder, @@ -753,6 +793,7 @@ class Coordinator: route_review=route_review, build_verify=build_verify, dispatch_node=dispatch_node_callable, + plan_gate=plan_gate, ) # The ResumeWorker is satisfied directly by the compiled LangGraph app @@ -1067,10 +1108,25 @@ class Coordinator: question = None if question is not None: - # Multi-turn: a new clarifier question is waiting. Post it to the - # transport (the drain path otherwise leaves it unposted) and tell - # the human more input is needed. Thread it (and its channel_ref) - # under the task's root message so the next reply maps back. + # A human gate is open. Two kinds (B2b): the plan-review DECISION + # gate (kind == 'plan_decision') and the clarifier QUESTION gate + # (kind == 'clarify' / legacy no-kind). The load-bearing signal is + # the pending-interrupt kind — NOT the task status, which still + # reads the review node's carried-over 'parked' while suspended at + # the gate (B2a). Branch on the kind so the plan gate posts its + # plan + findings presentation while the clarifier keeps its + # existing follow-up path. + if question.get("kind") == graph_mod.PLAN_DECISION_KIND: + self._post_plan_decision_gate( + question, label=label, root_ts=root_ts, short=short + ) + continue + + # Multi-turn clarify: a new clarifier question is waiting. Post it + # to the transport (the drain path otherwise leaves it unposted) + # and tell the human more input is needed. Thread it (and its + # channel_ref) under the task's root message so the next reply maps + # back. try: conn = connect(self._db_path) try: @@ -1160,6 +1216,234 @@ class Coordinator: thread_ts=root_ts, ) + # ------------------------------------------------------------------ # + # Plan-review decision gate (B2b). + # ------------------------------------------------------------------ # + + # Slack section-block text caps at ~3000 chars; the gate presentation budgets + # a bit under that for safety once the decision instructions are appended. + _GATE_PRESENTATION_BUDGET = 2700 + _GATE_TRUNCATION_NOTE = "\n…(truncated — full plan on the dashboard)" + + def _post_plan_decision_gate( + self, + question: "dict[str, Any]", + *, + label: str, + root_ts: str | None, + short: str, + ) -> None: + """Open + present the plan-review decision gate for a suspended task (B2b). + + The graph has suspended on the resumable PLAN_GATE interrupt (the + review-cap dead-end, B2a) whose payload mirrors the clarifier's PLUS + ``kind == 'plan_decision'``, ``plan``, and ``findings``. This: + + 1. Writes the durable ``pending_questions`` row with ``kind='plan_decision'`` + (24h deadline like the clarifier), threaded under the task root — via + :meth:`_open_plan_decision_row_if_absent` so the single-open-gate + invariant (B-2) holds (the clarifier row is already answered before the + plan stage runs, so a second open row would be a bug; we log + skip it). + 2. POSTs a presentation message — the plan (:meth:`_summarize_plan`) + the + review findings (:meth:`_summarize_blocker`) + the decision + instructions — threaded under the task root, truncated to Slack's block + limit. + + The decision answer then flows back through the UNCHANGED + ``submit_answer`` → resume-queue → ResumeWorker path (the graph's + ``plan_gate_node`` consumes the resume value); nothing here special-cases + the resume worker. Best-effort + fully guarded: a gate-post failure leaves + the durable interrupt in place (recovery re-derives it) and never breaks + the tick loop. + """ + deadline = question.get("deadline") or self._default_deadline() + thread_id = str(question.get("thread_id") or "") + question_id = str(question.get("question_id") or "") + turn = int(question.get("turn") or 0) + transport_name = str(question.get("transport") or "") + + # 1. Durable ledger row first (open, kind='plan_decision'), guarded by the + # single-open-gate invariant. + try: + conn = connect(self._db_path) + try: + opened = self._open_plan_decision_row_if_absent( + conn, + thread_id=thread_id, + question_id=question_id, + turn=turn, + transport_name=transport_name, + deadline=deadline, + root_ts=root_ts, + ) + finally: + conn.close() + except Exception: # noqa: BLE001 - a ledger error must not break the tick + _LOG.warning( + "failed to open plan-decision gate row for %s", short, exc_info=True + ) + opened = False + + if not opened: + # An open row already exists for this thread (single-open-gate + # invariant): do not present a second gate. Already logged in the + # opener; just stop here. + return + + # 2. Present the plan + findings + decision instructions, threaded under + # the task root. + body = self._plan_decision_presentation(question, label=label) + self._emit(body, thread_ts=root_ts) + + def _plan_decision_presentation( + self, question: "dict[str, Any]", *, label: str + ) -> str: + """Render the plan-gate presentation (plan + findings + instructions, B2b). + + Reuses :meth:`_summarize_plan` (plan shape) and :meth:`_summarize_blocker` + (review findings) so the gate view matches the rest of the lifecycle + threading, then appends the explicit decision instructions. The combined + body is truncated to Slack's section-block limit + (:data:`_GATE_PRESENTATION_BUDGET`) with a "(truncated — full plan on the + dashboard)" note so a multi-KB plan never overruns the block. + """ + plan_view = self._summarize_plan({"plan": question.get("plan")}) + findings = " ".join(str(question.get("findings") or "").split()) + if not findings: + findings = "(no review findings recorded)" + elif len(findings) > 600: + findings = findings[:600] + "…" + + instructions = ( + "• Decide: reply *approve* / *request changes * / *abandon* " + "in this thread." + ) + header = f"🧭 {label} — plan needs your decision (review could not approve it)." + + body = "\n".join( + [ + header, + plan_view, + f"• Review findings: {findings}", + instructions, + ] + ) + return self._truncate_for_slack(body) + + @classmethod + def _truncate_for_slack(cls, body: str) -> str: + """Truncate ``body`` to the gate presentation budget with a marker (B2b).""" + if len(body) <= cls._GATE_PRESENTATION_BUDGET: + return body + keep = cls._GATE_PRESENTATION_BUDGET - len(cls._GATE_TRUNCATION_NOTE) + return body[: max(keep, 0)].rstrip() + cls._GATE_TRUNCATION_NOTE + + def _open_plan_decision_row_if_absent( + self, + conn: Any, + *, + thread_id: str, + question_id: str, + turn: int, + transport_name: str, + deadline: str, + root_ts: str | None, + ) -> bool: + """Open a ``kind='plan_decision'`` ledger row, enforcing one-open-gate (B-2). + + The single-open-gate invariant: a thread has AT MOST ONE open + ``pending_questions`` row at a time. The clarifier row is already answered + before the plan stage runs, so under normal operation no open row exists + here. If one somehow does (a bug, or a redelivered drain re-presenting the + same gate), we do NOT open a second — we log + skip and return ``False`` + so the caller does not re-present. Returns ``True`` iff a fresh row was + opened + posted. + + The row is written ``open`` first (durable before the post), then the + transport posts the gate presentation threaded under ``root_ts`` and the + row's ``channel_ref`` is set to the root ts (so an inbound reply's + ``thread_ts`` maps back), mirroring :func:`responder.notify_question`. + """ + existing = _open_question_id_for_thread(conn, thread_id) + if existing is not None: + _LOG.warning( + "single-open-gate invariant: thread %s already has open question " + "%s; not opening a second (plan_decision) gate row", + thread_id[:8], + existing, + ) + return False + + # Durable row first (open, no ref), kind='plan_decision'. + from agent_team import ledger as ledger_mod # noqa: PLC0415 + + ledger_mod.post_question( + conn, + question_id=question_id, + thread_id=thread_id, + turn=turn, + transport=transport_name or type(self._transport).__name__, + deadline_at=deadline, + kind=graph_mod.PLAN_DECISION_KIND, + ) + + # Side-effecting post: thread the gate presentation under the root and + # record the channel_ref. A post failure is recoverable (row stays open, + # no ref) — swallow it exactly like notify_question. + channel_ref: str | None = None + try: + post_kwargs: dict[str, Any] = {} + if root_ts: + post_kwargs["thread_ts"] = root_ts + posted_ref = self._transport.post_question( + thread_id=thread_id, + question_id=question_id, + turn=turn, + question_set=self._plan_decision_question_set( + thread_id=thread_id, question_id=question_id, turn=turn + ), + deadline=deadline, + **post_kwargs, + ) + channel_ref = root_ts if root_ts else posted_ref + except Exception: # noqa: BLE001 - lost post is recoverable; keep the row + _LOG.warning( + "plan-decision gate post failed for thread %s; row stays open " + "for redelivery", + thread_id[:8], + exc_info=True, + ) + return True + + if channel_ref: + ledger_mod.set_channel_ref( + conn, question_id=question_id, channel_ref=channel_ref + ) + return True + + @staticmethod + def _plan_decision_question_set( + *, thread_id: str, question_id: str, turn: int + ) -> "QuestionSet": + """Build a minimal QuestionSet for the gate post (answer-mapping only). + + The transport's ``post_question`` requires a ``question_set`` so the + inbound answer can map back to ``question_id``; the gate's human-readable + presentation is posted via the lifecycle ``_emit`` sink, so this set only + needs to carry the identity (the decision-instruction wording lives in the + presentation message). The single question text is a terse decision + prompt for transports that render the set directly. + """ + from agent_team.transport.base import QuestionSet # noqa: PLC0415 + + return QuestionSet( + thread_id=thread_id, + question_id=question_id, + turn=turn, + questions=["Approve, request changes, or abandon this plan?"], + context={"kind": graph_mod.PLAN_DECISION_KIND}, + ) + @staticmethod def _summarize_plan(values: "dict[str, Any]") -> str: """Condensed, Slack-friendly view of the approved plan (summary + phases). @@ -1603,8 +1887,70 @@ class Coordinator: parked state for P1); this raises the injected ALARM hook so the stall is surfaced rather than silently spun on. Kept separate so the park policy is one obvious, testable place. + + **Plan-decision expiry (B2b).** A ``kind='plan_decision'`` gate row is an + ordinary ``pending_questions`` row, so the same deadline sweep expires it. + When the expired row is a plan-decision gate we ALSO emit a lifecycle + notice naming the task + that the plan DECISION expired unanswered + the + operator recovery path (re-assign / force-resume), threaded under the + task root, so an unanswered gate reads sensibly rather than as a generic + "clarifier question expired". Best-effort + guarded — never breaks tick. """ self._alarm_hook(question_id) + self._maybe_notify_plan_decision_expiry(question_id) + + def _maybe_notify_plan_decision_expiry(self, question_id: str) -> None: + """Emit a Slack lifecycle notice when a plan-decision gate row expires (B2b). + + Looks up the just-expired row; if it is a ``plan_decision`` gate it posts + a notice naming the task (read from the live graph state) + the recovery + path, threaded under the task root. A clarifier expiry is left to the + existing ALARM path (no extra notice). Fully guarded: any read/post error + is logged and swallowed so a notice failure never breaks the deadline + sweep. + """ + try: + conn = connect(self._db_path) + try: + from agent_team import ledger as ledger_mod # noqa: PLC0415 + + row = ledger_mod.get_question(conn, question_id) + finally: + conn.close() + except Exception: # noqa: BLE001 - a read error must not break the sweep + _LOG.warning( + "could not read expired question %s for plan-decision notice", + question_id, + exc_info=True, + ) + return + + if row is None or row.kind != graph_mod.PLAN_DECISION_KIND: + return + + # Name the task + thread the notice under its root, reading the live state. + label = f"(`{row.thread_id[:8]}`)" + root_ts: str | None = None + try: + snap = self._graph.get_state(graph_mod.thread_config(row.thread_id)) + values = getattr(snap, "values", {}) or {} + desc = str(values.get("task") or "").strip() + if desc: + if len(desc) > 90: + desc = desc[:90] + "…" + label = f'"{desc}" ({label})' + root_ts = str(values.get("slack_thread_ts") or "") or None + except Exception: # noqa: BLE001 - fall back to the bare id label + pass + + self._emit( + f"⌛ PLAN DECISION EXPIRED — {label}\n" + "• The plan-review decision was not answered within 24h, so the task " + "is parked.\n" + "• Recovery: re-assign the task, or force-resume it with a decision " + "(approve / request changes / abandon).", + thread_ts=root_ts, + ) @staticmethod def _default_alarm_hook(question_id: str) -> None: diff --git a/agent-team/tests/test_coordinator.py b/agent-team/tests/test_coordinator.py index 375036e..07b7d2a 100644 --- a/agent-team/tests/test_coordinator.py +++ b/agent-team/tests/test_coordinator.py @@ -1659,3 +1659,298 @@ def test_tick_ci_watch_parks_fail_closed_on_error(db_path: Path) -> None: assert report is not None assert report.parked_error == 1 assert alarms == ["t-err"] + + +# --------------------------------------------------------------------------- # +# Plan-review decision gate (Phase B2b) — coordinator wiring +# --------------------------------------------------------------------------- # + + +def _p2_plan_stub(state: Any) -> Any: + """Plan node that emits a plan and advances to REVIEW (no model call).""" + from agent_team.task_model import Phase, PipelineState, TaskStatus + + revisions = len(state.get("review_verdicts") or []) + return PipelineState( + plan={ + "summary": "ship the thing", + "phases": [{"name": "do it"}], + "revision": revisions, + }, + current_phase=Phase.REVIEW.value, + status=TaskStatus.ACTIVE.value, + ) + + +def _make_gate_coordinator( + db_path: Path, + *, + review_text: str = "VERDICT: REQUEST CHANGES\nstill not ready", + review_invoker: Any = None, + review_rounds: int | None = None, + transport: Transport | None = None, + notify: Any = None, + plan_node: Any = None, +) -> Coordinator: + """Build a gate-enabled Coordinator: stub clarify + real review loop + plan_gate. + + The review invoker NEVER approves (by default), so the plan<->review loop hits + the review cap and — with ``plan_gate=True`` — suspends on the resumable + PLAN_GATE interrupt instead of terminally parking. The clarify node stays the + deterministic stub (no Claude). + """ + from agent_team.nodes import review_loop + + saver = _Saver() + invoker = review_invoker or (lambda prompt, **kw: review_text) + review_loop.set_review_invoker(invoker) + + def _review_wiring() -> Any: + kw = {"max_review_rounds": review_rounds} if review_rounds is not None else {} + return review_loop.bind_review_node(kw), review_loop.route_after_review + + 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 _p2_plan_stub, + review_wiring=_review_wiring, + build_checkpointer=lambda _path: saver, + notify=notify, + plan_gate=True, + ) + + +@pytest.fixture() +def _restore_review_invoker(): + from agent_team.nodes import review_loop + + saved = review_loop._review_invoker + yield + review_loop._review_invoker = saved + + +def _drive_to_gate(coord: Coordinator, *, root_ts: str = "ROOT.TS") -> str: + """Start a task, answer the clarifier, and drain so it lands at the plan gate.""" + thread_id = coord.start_task( + task_text="add a thing", transport_name="slack", slack_thread_ts=root_ts + ) + qid = _only_open_row(db_path_of(coord))["question_id"] + coord.submit_answer({"question_id": qid, "answer": "scope is X", "via": "v"}) + results = coord.drain_resumes() + coord._post_resume_followups(results) + return thread_id + + +def db_path_of(coord: Coordinator) -> Path: + return coord._db_path + + +def test_setup_with_plan_gate_off_when_no_review_wiring(db_path: Path) -> None: + # A P1-only coordinator (no review wiring) must build cleanly even though the + # ctor flag defaults plan_gate True: setup forces it off (no review_node). + coord = _make_coordinator(db_path) # plan_gate defaults True, no review_wiring + coord.setup() + assert coord.graph is not None + + +def test_gate_writes_plan_decision_row_and_presents_plan( + db_path: Path, _restore_review_invoker: Any +) -> None: + posted: list[tuple[str, str | None]] = [] + coord = _make_gate_coordinator( + db_path, + notify=lambda message, thread_ts=None: posted.append((message, thread_ts)), + ) + coord.setup() + _drive_to_gate(coord, root_ts="ROOT.TS") + + # A durable open row with kind='plan_decision' exists. + row = _only_open_row(db_path) + assert row["kind"] == "plan_decision" + + # A presentation was posted, threaded under the task root, containing the + # plan summary + the review findings + the decision instructions. + gate = [p for p in posted if "plan needs your decision" in p[0]] + assert len(gate) == 1 + message, thread_ts = gate[0] + assert thread_ts == "ROOT.TS" + assert "ship the thing" in message # the plan summary is presented + assert "still not ready" in message # the review findings are presented + assert "approve" in message and "abandon" in message + + +def test_gate_approve_settles_task_not_interrupted( + db_path: Path, _restore_review_invoker: Any +) -> None: + coord = _make_gate_coordinator(db_path) + coord.setup() + thread_id = _drive_to_gate(coord) + + # Approve the gate via the SAME submit_answer -> resume path (unchanged). + qid = _only_open_row(db_path)["question_id"] + coord.submit_answer( + {"question_id": qid, "answer": {"decision": "approve"}, "via": "v"} + ) + coord.drain_resumes() + + # The task settled approved (phase BUILD, status ACTIVE) and is no longer + # interrupted. + assert graph_mod.pending_question(coord.graph, thread_id=thread_id) is None + state = graph_mod.get_pipeline_state(coord.graph, thread_id=thread_id) + assert state["current_phase"] == "build" + assert state["status"] == "active" + + +def test_gate_request_changes_reenters_pipeline( + db_path: Path, _restore_review_invoker: Any +) -> None: + # First review round REQUEST CHANGES (reach gate); after request_changes the + # next review APPROVES so the loop re-enters plan->review and settles. + texts = iter(["VERDICT: REQUEST CHANGES\nnope", "VERDICT: APPROVE\nnow good"]) + last = {"v": "VERDICT: APPROVE\nnow good"} + + def invoker(prompt: Any, **kw: Any) -> str: + try: + last["v"] = next(texts) + except StopIteration: + pass + return last["v"] + + coord = _make_gate_coordinator(db_path, review_invoker=invoker, review_rounds=1) + coord.setup() + thread_id = _drive_to_gate(coord) + + qid = _only_open_row(db_path)["question_id"] + coord.submit_answer( + { + "question_id": qid, + "answer": {"decision": "request_changes", "notes": "do X"}, + "via": "v", + } + ) + coord.drain_resumes() + + # The task re-entered the pipeline and settled (approved on the re-plan). + state = graph_mod.get_pipeline_state(coord.graph, thread_id=thread_id) + assert state["current_phase"] == "build" + # The human notes were folded in as a synthetic verdict. + assert any(v.get("reviewer") == "human_plan_gate" for v in state["review_verdicts"]) + + +def test_gate_abandon_fails_task(db_path: Path, _restore_review_invoker: Any) -> None: + coord = _make_gate_coordinator(db_path) + coord.setup() + thread_id = _drive_to_gate(coord) + + qid = _only_open_row(db_path)["question_id"] + coord.submit_answer( + {"question_id": qid, "answer": {"decision": "abandon"}, "via": "v"} + ) + coord.drain_resumes() + + state = graph_mod.get_pipeline_state(coord.graph, thread_id=thread_id) + assert state["status"] == "failed" + + +def test_single_open_gate_invariant_skips_second_open( + db_path: Path, _restore_review_invoker: Any +) -> None: + # Re-presenting the gate (a redelivered drain) must NOT open a second row. + posted: list[tuple[str, str | None]] = [] + coord = _make_gate_coordinator( + db_path, + notify=lambda message, thread_ts=None: posted.append((message, thread_ts)), + ) + coord.setup() + thread_id = _drive_to_gate(coord) + + open_rows_before = _all_rows(db_path, status="open") + assert len(open_rows_before) == 1 + + # Force a second presentation of the same gate (simulating a redelivered + # resume settling on the same interrupt). The single-open-gate guard must + # reject the second open. + question = graph_mod.pending_question(coord.graph, thread_id=thread_id) + coord._post_plan_decision_gate( + question, label='"x" (`short`)', root_ts="ROOT.TS", short="short123" + ) + + open_rows_after = _all_rows(db_path, status="open") + assert len(open_rows_after) == 1 # still exactly one open gate row + # No second presentation was posted. + gate = [p for p in posted if "plan needs your decision" in p[0]] + assert len(gate) == 1 + + +def test_clarifier_followup_still_posts_no_regression(db_path: Path) -> None: + # A clarifier-kind interrupt keeps the existing follow-up behavior (the gate + # branch must not steal it). Reuses the legacy multi-turn follow-up shape. + from agent_team import coordinator as coord_mod + + posted: list[Any] = [] + msgs: list[str] = [] + coord = _make_coordinator(db_path, notify=lambda m, **k: msgs.append(m)) + coord.setup() + coord.start_task(task_text="x", transport_name="slack") + + # A pending CLARIFY interrupt (no 'kind' => clarify); a stub notify_question. + coord._notify = lambda m, **k: msgs.append(m) + + class _QS: + thread_id = "abc12345deadbeef" + + import unittest.mock as _mock + + with ( + _mock.patch.object( + coord_mod.graph_mod, + "pending_question", + lambda _g, *, thread_id: {"question_set": _QS(), "deadline": "2099-01-01"}, + ), + _mock.patch.object( + coord_mod.responder_mod, + "notify_question", + lambda conn, transport, qset, *, deadline, thread_ts=None: posted.append( + qset + ), + ), + ): + coord._post_resume_followups([_resume_result("abc12345deadbeef")]) + + assert len(posted) == 1 # the clarifier follow-up was delivered + assert any("needs more input" in m for m in msgs) + + +def test_plan_decision_expiry_posts_recovery_notice( + db_path: Path, _restore_review_invoker: Any +) -> None: + # When a plan-decision gate row expires, tick() posts a recovery notice + # naming the task + the recovery path (re-assign / force-resume). + posted: list[tuple[str, str | None]] = [] + coord = _make_gate_coordinator( + db_path, + notify=lambda message, thread_ts=None: posted.append((message, thread_ts)), + ) + coord.setup() + _drive_to_gate(coord, root_ts="ROOT.TS") + posted.clear() + + # Force the open gate row's deadline into the past, then sweep. + conn = connect(db_path) + try: + conn.execute( + "UPDATE pending_questions SET deadline_at=? WHERE status='open'", + ("2000-01-01T00:00:00+00:00",), + ) + conn.commit() + finally: + conn.close() + + coord.tick() + + notice = [p for p in posted if "PLAN DECISION EXPIRED" in p[0]] + assert len(notice) == 1 + message, thread_ts = notice[0] + assert thread_ts == "ROOT.TS" + assert "re-assign" in message.lower() or "force-resume" in message.lower()