diff --git a/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index 312ff62..16af0ea 100644 --- a/agent-team/agent_team/coordinator.py +++ b/agent-team/agent_team/coordinator.py @@ -450,6 +450,7 @@ class Coordinator: alarm_hook: AlarmHook | None = None, build_listener: ListenerFactory | None = None, new_task_callback: "Callable[[str, str], str] | None" = None, + notify: "Callable[[str], None] | None" = None, ) -> None: self._db_path = Path(db_path) self._transport = transport @@ -482,6 +483,11 @@ class Coordinator: # None, the listener ignores /new-task. The serve path wires the # coordinator's own start_task adapter via set_new_task_callback(). self._new_task_callback = new_task_callback + # Optional lifecycle-notification sink (posts a plain status line to the + # operator channel, e.g. Slack #agent-team). Left None = silent (existing + # behavior). The serve path injects a Slack poster so a task is never a + # black box: the human sees parked / needs-more-input / plan-ready. + self._notify = notify # Built by setup(). self._graph: Any = None @@ -741,6 +747,81 @@ class Coordinator: # Maintenance tick (deadline policy + drain). # ------------------------------------------------------------------ # + def _emit(self, message: str) -> None: + """Post a lifecycle status line to the notify sink (never raises). + + A notification failure must never disturb the pipeline, so the sink call + is wrapped; a broken Slack post is logged and swallowed. + """ + if self._notify is None: + return + try: + self._notify(message) + except Exception: # noqa: BLE001 - a notify failure must not break the loop + _LOG.warning("notify sink raised; dropping status message", exc_info=True) + + def _post_resume_followups(self, results: list[ResumeResult]) -> None: + """After a drain, deliver follow-up questions + emit lifecycle milestones. + + For each resumed thread, the graph has settled at one of: a NEW pending + question (multi-turn clarify — which the normal drain path does NOT post, + only the startup recover sweep did, so post it here), a PARKED terminal + state, or a completed/plan-ready terminal state. Emits a human-readable + status line for each so an answered task is never a black box. Best-effort + and fully guarded — a follow-up failure never breaks the tick loop. + """ + from agent_team.task_model import TaskStatus # noqa: PLC0415 + + seen: set[str] = set() + for result in results: + thread_id = getattr(result, "thread_id", None) + if not thread_id or thread_id in seen: + continue + seen.add(thread_id) + short = thread_id[:8] + try: + question = graph_mod.pending_question(self._graph, thread_id=thread_id) + except Exception: # noqa: BLE001 - never let a status check break the loop + 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. + try: + conn = connect(self._db_path) + try: + responder_mod.notify_question( + conn, + self._transport, + question["question_set"], + deadline=question.get("deadline") + or self._default_deadline(), + ) + finally: + conn.close() + except Exception: # noqa: BLE001 - post failure must not break tick + _LOG.warning( + "failed to post follow-up question for %s", short, exc_info=True + ) + self._emit(f"❓ Task {short}: needs more input — posted a question.") + continue + + # No pending question: the task settled. Distinguish parked vs done. + try: + snap = self._graph.get_state(graph_mod.thread_config(thread_id)) + values = getattr(snap, "values", {}) or {} + status = values.get("status") + except Exception: # noqa: BLE001 + status = None + if status == TaskStatus.PARKED.value: + self._emit( + f"⚠️ Task {short}: parked — needs your attention " + "(plan/review escalation). Re-assign or steer it to resume." + ) + else: + self._emit(f"✅ Task {short}: plan ready for review.") + def tick(self) -> list[ResumeResult]: """One maintenance pass: deadline sweep + park policy, then drain (§3.3.1). @@ -760,7 +841,9 @@ class Coordinator: for question_id in expired: self._park(question_id) - return self.drain_resumes() + results = self.drain_resumes() + self._post_resume_followups(results) + return results def _park(self, question_id: str) -> None: """Apply the park policy to one expired question (§6.6 ALARM, not spin). diff --git a/agent-team/run-team.py b/agent-team/run-team.py index 9248ccd..57f9b09 100644 --- a/agent-team/run-team.py +++ b/agent-team/run-team.py @@ -61,6 +61,7 @@ from __future__ import annotations import argparse import getpass import json +import logging import os import sqlite3 import sys @@ -106,6 +107,7 @@ _DEFAULT_DB = _CLI_DIR / "state" / "agent_team.sqlite" # Default audit log for destructive actions, alongside the ledger DB. _DEFAULT_AUDIT_LOG = _CLI_DIR / "state" / "audit.log.jsonl" +_LOG = logging.getLogger("agent_team.run_team") # Columns selected for list/show rendering, in display order. _QUESTION_COLUMNS: tuple[str, ...] = ( @@ -514,6 +516,50 @@ def _build_context_provider() -> "Callable[[], str]": return load_handbook_conventions +def _build_notifiers( + args: argparse.Namespace, +) -> "tuple[Callable[[str], None] | None, Callable[[str], None] | None]": + """Build the (notify, alarm_hook) Slack notifiers for the coordinator. + + Returns ``(None, None)`` for dry-run / non-Slack / no-channel so import, + ``--help``, ledger commands, and token-less dry runs stay silent and need no + Slack credentials. For live Slack with ``SLACK_CHANNEL_ID`` set, ``notify`` + posts a plain status line to the channel (via the same ``build_slack_poster`` + the transport uses), and ``alarm_hook`` logs the deadline-park WARNING AND + posts a parked-task notice. Both are best-effort — the coordinator wraps the + notify sink so a Slack failure never disturbs the pipeline. + """ + if getattr(args, "dry_run", False) or args.transport != "slack": + return None, None + channel = os.environ.get("SLACK_CHANNEL_ID", "") + if not channel: + return None, None + try: + from agent_team.transport.slack_live import build_slack_poster + + poster = build_slack_poster() + except Exception: # noqa: BLE001 - no token / SDK -> run without notifications + _LOG.warning( + "Slack notifier unavailable; coordinator runs without notifications" + ) + return None, None + + def notify(message: str) -> None: + poster({"channel": channel, "text": message}) + + def alarm_hook(question_id: str) -> None: + _LOG.warning("park ALARM: clarifier question %s expired", question_id) + try: + notify( + f"⚠️ Task parked: clarifier question {question_id[:8]} expired with " + "no answer in the window. Re-assign or answer to resume." + ) + except Exception: # noqa: BLE001 - notify failure must not break the park path + pass + + return notify, alarm_hook + + def _build_coordinator(args: argparse.Namespace) -> Any: """Construct a :class:`Coordinator` for the ``start`` / ``serve`` commands. @@ -543,6 +589,12 @@ def _build_coordinator(args: argparse.Namespace) -> Any: # unavailable. context_provider = _build_context_provider() + # Lifecycle notifications: post plain status lines to the Slack channel so a + # task is never a black box (parked / needs-more-input / plan-ready, and + # deadline-park ALARMs). Live-Slack only; dry-run / non-Slack / no-channel = + # silent (notify None) so import + ledger commands need no token. + notify, alarm_hook = _build_notifiers(args) + # Production runs the full P2 graph: the wrapped real planner + the bound # GPT-4.1 review loop (Plane-2 depth-first). These factories are lazy and # only build/bind the model seams when a task actually runs. @@ -556,6 +608,8 @@ def _build_coordinator(args: argparse.Namespace) -> Any: context_provider=context_provider ), review_wiring=default_review_wiring, + notify=notify, + alarm_hook=alarm_hook, ) # WS2: an allowlisted Slack /new-task starts a task on THIS coordinator. Set # post-construction (the adapter closes over the just-built coordinator), and diff --git a/agent-team/tests/test_coordinator.py b/agent-team/tests/test_coordinator.py index 0692e35..b305552 100644 --- a/agent-team/tests/test_coordinator.py +++ b/agent-team/tests/test_coordinator.py @@ -912,3 +912,101 @@ def test_serve_wires_start_and_stop_around_the_loop( assert listener.served is True # started before the loop assert listener.closed is True # stopped in finally on interrupt + + +# --------------------------------------------------------------------------- # +# Lifecycle notifications (_post_resume_followups + _emit): an answered task is +# never a black box — parked / needs-more-input / plan-ready post to the sink, +# and a follow-up clarifier question is delivered (the drain path otherwise +# leaves multi-turn questions unposted). +# --------------------------------------------------------------------------- # + + +def _resume_result(thread_id: str) -> Any: + from agent_team.resume_worker import ResumeOutcome, ResumeResult + + return ResumeResult( + outcome=ResumeOutcome.RESUMED, + thread_id=thread_id, + question_id="q-1", + turn=1, + graph_result=None, + ) + + +def test_followups_posts_new_question_and_emits_needs_input( + db_path: Path, monkeypatch: Any +) -> None: + from agent_team import coordinator as coord_mod + + posted: list[Any] = [] + msgs: list[str] = [] + coord = _make_coordinator(db_path) + coord._notify = msgs.append + coord.setup() + + monkeypatch.setattr( + coord_mod.graph_mod, + "pending_question", + lambda _g, *, thread_id: {"question_set": object(), "deadline": "2099-01-01"}, + ) + monkeypatch.setattr( + coord_mod.responder_mod, + "notify_question", + lambda conn, transport, qset, *, deadline: posted.append((qset, deadline)), + ) + + coord._post_resume_followups([_resume_result("abc12345deadbeef")]) + assert len(posted) == 1 # the follow-up question was delivered + assert any("needs more input" in m for m in msgs) + assert any("abc12345" in m for m in msgs) + + +def test_followups_emits_parked(db_path: Path, monkeypatch: Any) -> None: + from agent_team import coordinator as coord_mod + from agent_team.task_model import TaskStatus + + msgs: list[str] = [] + coord = _make_coordinator(db_path) + coord._notify = msgs.append + coord.setup() + + monkeypatch.setattr( + coord_mod.graph_mod, "pending_question", lambda _g, *, thread_id: None + ) + + class _Snap: + values = {"status": TaskStatus.PARKED.value} + + monkeypatch.setattr(coord._graph, "get_state", lambda _cfg: _Snap()) + coord._post_resume_followups([_resume_result("dead0001beef")]) + assert any("parked" in m.lower() for m in msgs) + + +def test_followups_emits_plan_ready(db_path: Path, monkeypatch: Any) -> None: + from agent_team import coordinator as coord_mod + + msgs: list[str] = [] + coord = _make_coordinator(db_path) + coord._notify = msgs.append + coord.setup() + + monkeypatch.setattr( + coord_mod.graph_mod, "pending_question", lambda _g, *, thread_id: None + ) + + class _Snap: + values = {"status": "active"} + + monkeypatch.setattr(coord._graph, "get_state", lambda _cfg: _Snap()) + coord._post_resume_followups([_resume_result("feed0002face")]) + assert any("plan ready" in m.lower() for m in msgs) + + +def test_emit_swallows_notify_failure(db_path: Path) -> None: + def _boom(_msg: str) -> None: + raise RuntimeError("slack down") + + coord = _make_coordinator(db_path) + coord._notify = _boom + coord._emit("anything") # must not raise diff --git a/agent-team/tests/test_run_team.py b/agent-team/tests/test_run_team.py index 9d180ae..2826bc5 100644 --- a/agent-team/tests/test_run_team.py +++ b/agent-team/tests/test_run_team.py @@ -647,9 +647,13 @@ class _FakeCoordinator: build_clarify_node: Any = None, build_plan_node: Any = None, review_wiring: Any = None, + notify: Any = None, + alarm_hook: Any = None, ) -> None: self.db_path = db_path self.transport = transport + self.notify = notify + self.alarm_hook = alarm_hook # The production CLI opts the coordinator into the P2 graph by injecting # these factories; record them so the wiring is asserted, not ignored. self.build_clarify_node = build_clarify_node