feat(notify): Slack lifecycle notifications + deliver multi-turn questions

The bot only ever posted clarifier questions; parks/completions were silent and
multi-turn follow-up questions were never posted during normal operation (only
the startup recover sweep posted them). So an answered task was a black box.

- Coordinator gains a 'notify' sink + _emit() (guarded, never breaks the loop).
- tick() now runs _post_resume_followups(results) after the drain: for each
  resumed thread it (a) POSTS a newly-pending clarifier question (fixes silent
  multi-turn — the drain path left it unposted) and emits 'needs more input';
  else emits 'parked — needs attention' or 'plan ready for review' from the
  settled state.
- run-team serve wires notify -> Slack channel (build_slack_poster) and an
  alarm_hook that logs the deadline-park WARNING AND posts a parked notice.
  Live-Slack only; dry-run/non-Slack/no-channel = silent (None), no token needed.

Tests: +4 (needs-input/parked/plan-ready emits + notify-failure swallow);
_FakeCoordinator gains notify/alarm_hook. 1147 passed, ruff clean.
This commit is contained in:
Adam Moussa 2026-06-23 14:17:08 -04:00
parent a6fd1cfc75
commit f0a2dc27a5
4 changed files with 240 additions and 1 deletions

View file

@ -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).

View file

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

View file

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

View file

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