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:
parent
6d76194474
commit
ebc10649a2
4 changed files with 240 additions and 1 deletions
|
|
@ -450,6 +450,7 @@ class Coordinator:
|
||||||
alarm_hook: AlarmHook | None = None,
|
alarm_hook: AlarmHook | None = None,
|
||||||
build_listener: ListenerFactory | None = None,
|
build_listener: ListenerFactory | None = None,
|
||||||
new_task_callback: "Callable[[str, str], str] | None" = None,
|
new_task_callback: "Callable[[str, str], str] | None" = None,
|
||||||
|
notify: "Callable[[str], None] | None" = None,
|
||||||
) -> None:
|
) -> None:
|
||||||
self._db_path = Path(db_path)
|
self._db_path = Path(db_path)
|
||||||
self._transport = transport
|
self._transport = transport
|
||||||
|
|
@ -482,6 +483,11 @@ class Coordinator:
|
||||||
# None, the listener ignores /new-task. The serve path wires the
|
# None, the listener ignores /new-task. The serve path wires the
|
||||||
# coordinator's own start_task adapter via set_new_task_callback().
|
# coordinator's own start_task adapter via set_new_task_callback().
|
||||||
self._new_task_callback = 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().
|
# Built by setup().
|
||||||
self._graph: Any = None
|
self._graph: Any = None
|
||||||
|
|
@ -741,6 +747,81 @@ class Coordinator:
|
||||||
# Maintenance tick (deadline policy + drain).
|
# 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]:
|
def tick(self) -> list[ResumeResult]:
|
||||||
"""One maintenance pass: deadline sweep + park policy, then drain (§3.3.1).
|
"""One maintenance pass: deadline sweep + park policy, then drain (§3.3.1).
|
||||||
|
|
||||||
|
|
@ -760,7 +841,9 @@ class Coordinator:
|
||||||
for question_id in expired:
|
for question_id in expired:
|
||||||
self._park(question_id)
|
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:
|
def _park(self, question_id: str) -> None:
|
||||||
"""Apply the park policy to one expired question (§6.6 ALARM, not spin).
|
"""Apply the park policy to one expired question (§6.6 ALARM, not spin).
|
||||||
|
|
|
||||||
|
|
@ -61,6 +61,7 @@ from __future__ import annotations
|
||||||
import argparse
|
import argparse
|
||||||
import getpass
|
import getpass
|
||||||
import json
|
import json
|
||||||
|
import logging
|
||||||
import os
|
import os
|
||||||
import sqlite3
|
import sqlite3
|
||||||
import sys
|
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 for destructive actions, alongside the ledger DB.
|
||||||
_DEFAULT_AUDIT_LOG = _CLI_DIR / "state" / "audit.log.jsonl"
|
_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.
|
# Columns selected for list/show rendering, in display order.
|
||||||
_QUESTION_COLUMNS: tuple[str, ...] = (
|
_QUESTION_COLUMNS: tuple[str, ...] = (
|
||||||
|
|
@ -514,6 +516,50 @@ def _build_context_provider() -> "Callable[[], str]":
|
||||||
return load_handbook_conventions
|
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:
|
def _build_coordinator(args: argparse.Namespace) -> Any:
|
||||||
"""Construct a :class:`Coordinator` for the ``start`` / ``serve`` commands.
|
"""Construct a :class:`Coordinator` for the ``start`` / ``serve`` commands.
|
||||||
|
|
||||||
|
|
@ -543,6 +589,12 @@ def _build_coordinator(args: argparse.Namespace) -> Any:
|
||||||
# unavailable.
|
# unavailable.
|
||||||
context_provider = _build_context_provider()
|
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
|
# 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
|
# 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.
|
# 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
|
context_provider=context_provider
|
||||||
),
|
),
|
||||||
review_wiring=default_review_wiring,
|
review_wiring=default_review_wiring,
|
||||||
|
notify=notify,
|
||||||
|
alarm_hook=alarm_hook,
|
||||||
)
|
)
|
||||||
# WS2: an allowlisted Slack /new-task starts a task on THIS coordinator. Set
|
# WS2: an allowlisted Slack /new-task starts a task on THIS coordinator. Set
|
||||||
# post-construction (the adapter closes over the just-built coordinator), and
|
# post-construction (the adapter closes over the just-built coordinator), and
|
||||||
|
|
|
||||||
|
|
@ -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.served is True # started before the loop
|
||||||
assert listener.closed is True # stopped in finally on interrupt
|
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
|
||||||
|
|
|
||||||
|
|
@ -647,9 +647,13 @@ class _FakeCoordinator:
|
||||||
build_clarify_node: Any = None,
|
build_clarify_node: Any = None,
|
||||||
build_plan_node: Any = None,
|
build_plan_node: Any = None,
|
||||||
review_wiring: Any = None,
|
review_wiring: Any = None,
|
||||||
|
notify: Any = None,
|
||||||
|
alarm_hook: Any = None,
|
||||||
) -> None:
|
) -> None:
|
||||||
self.db_path = db_path
|
self.db_path = db_path
|
||||||
self.transport = transport
|
self.transport = transport
|
||||||
|
self.notify = notify
|
||||||
|
self.alarm_hook = alarm_hook
|
||||||
# The production CLI opts the coordinator into the P2 graph by injecting
|
# The production CLI opts the coordinator into the P2 graph by injecting
|
||||||
# these factories; record them so the wiring is asserted, not ignored.
|
# these factories; record them so the wiring is asserted, not ignored.
|
||||||
self.build_clarify_node = build_clarify_node
|
self.build_clarify_node = build_clarify_node
|
||||||
|
|
|
||||||
Reference in a new issue