Merge pull request #54 from Sea-Haven-Industries/fix/agent-team-plan-presentation-threading

fix(agent-team): thread lifecycle milestones + present the plan in Slack
This commit is contained in:
Adam Moussa 2026-06-23 16:44:53 -04:00 • committed by GitHub
commit 778737fa9f
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
4 changed files with 162 additions and 3 deletions

View file

@ -995,11 +995,50 @@ class Coordinator:
thread_ts=root_ts,
)
else:
# Plan approved + settled at the P2 terminus. PRESENT the plan
# (condensed) so the human can actually review it in-thread, not
# just a "ready" notice with nothing to look at.
self._emit(
f"✅ {label} — plan ready for review (phase: {phase}).",
f"✅ {label} — plan ready for review (phase: {phase}).\n"
f"{self._summarize_plan(values)}",
thread_ts=root_ts,
)
@staticmethod
def _summarize_plan(values: "dict[str, Any]") -> str:
"""Condensed, Slack-friendly view of the approved plan (summary + phases).
Posts the plan ``summary`` (collapsed + truncated) plus the numbered
phase NAMES — enough to review/approve the shape in-thread without
dumping every step. Falls back to a bare line if the plan is malformed
(the milestone must still post). Step detail lives on the status
dashboard / a follow-up; this is the at-a-glance review view.
"""
plan = values.get("plan")
if not isinstance(plan, dict):
return "• (plan unavailable to summarize)"
lines: list[str] = []
summary = " ".join(str(plan.get("summary") or "").split())
if summary:
if len(summary) > 350:
summary = summary[:350] + "…"
lines.append(f"• Summary: {summary}")
phases = plan.get("phases") or []
if isinstance(phases, list) and phases:
names = [
str(p.get("name")).strip()
for p in phases
if isinstance(p, dict) and str(p.get("name") or "").strip()
]
if names:
lines.append(f"• Phases ({len(names)}):")
lines += [f" {i}. {n}" for i, n in enumerate(names, start=1)]
lines.append(
"• Reply in this thread to steer, or re-assign with changes. "
"(Full step detail: status dashboard.)"
)
return "\n".join(lines)
@staticmethod
def _summarize_blocker(values: "dict[str, Any]") -> str:
"""Human-readable reason a task parked, from the last review verdict.

View file

@ -544,8 +544,17 @@ def _build_notifiers(
)
return None, None
def notify(message: str) -> None:
poster({"channel": channel, "text": message})
def notify(message: str, thread_ts: str | None = None) -> None:
# `thread_ts` is the one-thread-per-task root ts: the coordinator's
# `_emit` passes it so lifecycle milestones (plan-ready / parked /
# failed) thread under the task's "📥 Task received" root. Without this
# param `_emit` hit a TypeError and silently fell back to a TOP-LEVEL
# post, so every milestone landed unthreaded. build_slack_poster already
# forwards a `thread_ts` payload key to chat.postMessage.
payload = {"channel": channel, "text": message}
if thread_ts:
payload["thread_ts"] = thread_ts
poster(payload)
def alarm_hook(question_id: str) -> None:
_LOG.warning("park ALARM: clarifier question %s expired", question_id)

View file

@ -409,6 +409,73 @@ def test_drain_resumes_one_failing_task_does_not_block_others(db_path: Path) ->
assert graph_mod.get_pipeline_state(coord.graph, thread_id=t2)["status"] == "failed"
# --------------------------------------------------------------------------- #
# plan presentation + lifecycle-milestone threading
# --------------------------------------------------------------------------- #
def test_summarize_plan_condensed_summary_and_phase_names() -> None:
"""The condensed plan view renders the summary + numbered phase NAMES only."""
values = {
"plan": {
"summary": "Add a hermetic smoke test.",
"phases": [
{"name": "Audit infra", "steps": ["look at tests/", "read conftest"]},
{
"name": "Write test_smoke.py",
"steps": ["import checks", "graph build"],
},
{"name": "Lint + commit", "steps": ["ruff", "pytest"]},
],
}
}
out = Coordinator._summarize_plan(values)
assert "Summary: Add a hermetic smoke test." in out
assert "Phases (3):" in out
assert "1. Audit infra" in out
assert "2. Write test_smoke.py" in out
assert "3. Lint + commit" in out
# Condensed: step detail is NOT dumped.
assert "read conftest" not in out
assert "ruff" not in out
def test_summarize_plan_malformed_falls_back_without_crashing() -> None:
assert "unavailable" in Coordinator._summarize_plan({"plan": None})
assert "unavailable" in Coordinator._summarize_plan({})
def test_plan_ready_milestone_threads_under_root_and_presents_plan(
db_path: Path,
) -> None:
"""A settled approved plan posts the plan (condensed) THREADED under the root.
Covers both fixes: the milestone is delivered with the task's
``slack_thread_ts`` (threading) and now carries the plan summary, not just a
bare "ready" line (presentation).
"""
posted: list[tuple[str, str | None]] = []
coord = _make_coordinator(
db_path,
notify=lambda message, thread_ts=None: posted.append((message, thread_ts)),
)
coord.setup()
coord.start_task(
task_text="add a smoke test", transport_name="slack", slack_thread_ts="ROOT.TS"
)
qid = _only_open_row(db_path)["question_id"]
coord.submit_answer({"question_id": qid, "answer": "go", "via": "v"})
results = coord.drain_resumes()
assert results[0].outcome is ResumeOutcome.RESUMED
coord._post_resume_followups(results)
ready = [p for p in posted if "plan ready for review" in p[0]]
assert len(ready) == 1
message, thread_ts = ready[0]
assert thread_ts == "ROOT.TS" # threaded under the task root, not top-level
assert "Summary:" in message # the plan is actually presented
# --------------------------------------------------------------------------- #
# tick — deadline sweep + park ALARM + drain
# --------------------------------------------------------------------------- #

View file

@ -1033,3 +1033,47 @@ def test_fix_non_fixable_finding_returns_one(
)
assert rc == 1
assert "FIX NOT PLANNED" in out.getvalue()
# --------------------------------------------------------------------------- #
# _build_notifiers — the Slack lifecycle-milestone sink (one-thread-per-task)
# --------------------------------------------------------------------------- #
def test_notify_sink_forwards_thread_ts(
cli: ModuleType, monkeypatch: pytest.MonkeyPatch
) -> None:
"""The notify sink must pass `thread_ts` so milestones thread under the task.
Regression: the sink was `def notify(message)` with no `thread_ts`, so the
coordinator's `_emit(message, thread_ts=root)` hit a TypeError and silently
fell back to a TOP-LEVEL post — every plan-ready/parked/failed milestone
landed unthreaded. build_slack_poster already forwards a `thread_ts` payload
key, so the only gap was this wrapper.
"""
captured: list[dict] = []
monkeypatch.setenv("SLACK_CHANNEL_ID", "C123")
# _build_notifiers imports build_slack_poster from slack_live at call time.
import agent_team.transport.slack_live as slack_live
monkeypatch.setattr(
slack_live,
"build_slack_poster",
lambda: lambda payload: captured.append(payload),
)
args = argparse.Namespace(dry_run=False, transport="slack")
notify, _alarm = cli._build_notifiers(args)
assert notify is not None
notify("threaded milestone", thread_ts="ROOT.TS")
notify("top-level milestone")
assert captured[0] == {
"channel": "C123",
"text": "threaded milestone",
"thread_ts": "ROOT.TS",
}
# No thread_ts when none is given (top-level post, not a broken key).
assert "thread_ts" not in captured[1]
assert captured[1] == {"channel": "C123", "text": "top-level milestone"}