From a96a5b487afe878ce451d28672593c84457634c9 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 12:38:43 -0400 Subject: [PATCH 01/13] feat(integration): wire WS5 context_provider + WS2 /new-task into serve Integration branch combining WS0-WS5 (PRs #43-#46) + the activation wiring that flips the safe seams ON in the run-team serve path: - WS5 (D10): inject the Sea Haven handbook conventions into the planner prompt via context_provider (zero-arg handbook loader; fail-safe to '' when absent). - WS2: an allowlisted Slack /new-task starts a task on this coordinator (set_new_task_callback adapter -> start_task; AUTHZ-01 gates it upstream). - WS1 bind_multi_invoker() is already wired in _cmd_serve. Coordinator gains a new_task_callback param + set_new_task_callback() (resolves the constructor chicken-and-egg of referencing the coordinator's own start_task); default_slack_listener_factory forwards it to the SlackListener. NOT wired (deliberately): the P3 dispatch_node / build_verify path. Activating it correctly needs a per-task expected_run_id bound into gated_build_verify_wiring (plumbing that does not exist yet) AND the CI trust-boundary security re-review. It stays inert pending that work. Tests: +5 activation-wiring tests; _FakeCoordinator stub gains set_new_task_callback. Full agent-team suite: 1140 passed, ruff clean. --- agent-team/agent_team/coordinator.py | 26 +++ agent-team/run-team.py | 47 ++++- agent-team/tests/test_run_team.py | 4 + agent-team/tests/test_ws_activation_wiring.py | 160 ++++++++++++++++++ 4 files changed, 234 insertions(+), 3 deletions(-) create mode 100644 agent-team/tests/test_ws_activation_wiring.py diff --git a/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index 623899d..a8336c6 100644 --- a/agent-team/agent_team/coordinator.py +++ b/agent-team/agent_team/coordinator.py @@ -152,6 +152,7 @@ def default_slack_listener_factory( transport: SlackTransport, db_path: Path, enqueue_resume: Callable[[Any], None], + new_task_callback: "Callable[[str, str], str] | None" = None, ) -> Any: """Build the live :class:`SlackListener` from the coordinator's seams (D-1). @@ -166,6 +167,11 @@ def default_slack_listener_factory( (AUTHZ-01), so this factory deliberately does not weaken that — it injects no ``owner_ids`` and lets ``serve`` read + enforce them. + ``new_task_callback`` (WS2) is forwarded to the listener so an allowlisted + ``/new-task`` command starts a task. Left ``None``, the listener ignores + ``/new-task`` (its built-in default) — the ``/new-task`` path still runs + AFTER the AUTHZ-01 owner check regardless. + Imported lazily for the same import-hygiene reason as the clarifier / planner factories (the listener pulls the transport + responder leaves). """ @@ -177,6 +183,7 @@ def default_slack_listener_factory( enqueue_resume, app_token=os.environ.get("SLACK_APP_TOKEN") or None, bot_token=os.environ.get("SLACK_BOT_TOKEN") or None, + new_task_callback=new_task_callback, ) @@ -434,6 +441,7 @@ class Coordinator: deadline_window: timedelta | None = None, alarm_hook: AlarmHook | None = None, build_listener: ListenerFactory | None = None, + new_task_callback: "Callable[[str, str], str] | None" = None, ) -> None: self._db_path = Path(db_path) self._transport = transport @@ -462,6 +470,10 @@ class Coordinator: # the start/no-start/shutdown wiring with no live socket. Only consulted # by serve() when Slack is the live transport AND the app token is set. self._build_listener = build_listener + # WS2: optional /new-task handler forwarded to the Slack listener. Left + # 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 # Built by setup(). self._graph: Any = None @@ -478,6 +490,19 @@ class Coordinator: # Accessors (the shared queue is the slack_listener handoff seam). # ------------------------------------------------------------------ # + def set_new_task_callback( + self, callback: "Callable[[str, str], str] | None" + ) -> None: + """Wire the WS2 ``/new-task`` handler (call before :meth:`serve`). + + Avoids the constructor chicken-and-egg of referencing the coordinator's + own ``start_task`` at build time: the serve path constructs the + coordinator, then sets ``lambda text, source: self.start_task(...)``. Must + be set before :meth:`_maybe_start_slack_listener` runs (i.e. before + :meth:`serve`); it is read when the listener is built. + """ + self._new_task_callback = callback + @property def resume_queue(self) -> "queue.Queue[Any]": """The shared resume-job queue (slack_listener enqueues, drainer drains).""" @@ -1028,6 +1053,7 @@ class Coordinator: transport=self._transport, db_path=self._db_path, enqueue_resume=self._resume_queue.put, + new_task_callback=self._new_task_callback, ) self._listener = listener thread = threading.Thread( diff --git a/agent-team/run-team.py b/agent-team/run-team.py index 7d07e3a..33d914f 100644 --- a/agent-team/run-team.py +++ b/agent-team/run-team.py @@ -66,7 +66,7 @@ import sqlite3 import sys from datetime import datetime, timezone from pathlib import Path -from typing import Any, Sequence +from typing import Any, Callable, Sequence # ``run-team.py`` lives in ``agent-team/`` next to the importable ``agent_team`` # package. The hyphenated filename cannot itself be imported, so when run as a @@ -492,6 +492,28 @@ def _cmd_supersede(args: argparse.Namespace, *, out: Any) -> int: return 0 +def _build_context_provider() -> "Callable[[], str]": + """Return the WS5 (D10) context provider: the Sea Haven handbook conventions. + + The ``context_provider`` seam is zero-arg (``Callable[[], str]``), so it + supplies STATIC context — the handbook conventions loaded from + ``SEA_HAVEN_HANDBOOK_DIR`` (or the default dir) via + :func:`agent_team.nodes.handbook.load_handbook_conventions`. That loader is + itself fail-safe (caps file count/size; returns ``""`` and never raises when + the dir is absent/empty/unreadable), so wiring it is safe even on a box where + the handbook has not been synced yet. + + (Task-keyed memory retrieval is NOT wired here: the seam takes no task text, + so it cannot form a retrieval query — that would need a task-aware seam.) + + Imported lazily for the same import-hygiene reason as the coordinator + factories (keeps ``--help`` / ledger commands import-clean). + """ + from agent_team.nodes.handbook import load_handbook_conventions + + return load_handbook_conventions + + def _build_coordinator(args: argparse.Namespace) -> Any: """Construct a :class:`Coordinator` for the ``start`` / ``serve`` commands. @@ -512,15 +534,34 @@ def _build_coordinator(args: argparse.Namespace) -> Any: ) transport = _build_transport(args) + # WS5 (D10): inject the Sea Haven handbook conventions into the planner prompt + # via the context_provider seam. The provider is the zero-arg handbook loader, + # which is itself fail-safe (returns "" when the handbook dir is absent), and + # planner.plan_node additionally swallows provider errors — so this never + # affects a run where the handbook is unavailable. + context_provider = _build_context_provider() + # 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. - return Coordinator( + coordinator = Coordinator( db_path=args.db, transport=transport, - build_plan_node=default_plan_node_factory, + build_plan_node=lambda: default_plan_node_factory( + context_provider=context_provider + ), review_wiring=default_review_wiring, ) + # WS2: an allowlisted Slack /new-task starts a task on THIS coordinator. Set + # post-construction (the adapter closes over the just-built coordinator), and + # before serve() builds the listener. AUTHZ-01 (owner allowlist) gates this + # upstream in the listener; the source label is always "slack". + coordinator.set_new_task_callback( + lambda task_text, _source: coordinator.start_task( + task_text=task_text, transport_name="slack" + ) + ) + return coordinator def _build_transport(args: argparse.Namespace) -> Any: diff --git a/agent-team/tests/test_run_team.py b/agent-team/tests/test_run_team.py index 86b3574..7ca0ce7 100644 --- a/agent-team/tests/test_run_team.py +++ b/agent-team/tests/test_run_team.py @@ -655,11 +655,15 @@ class _FakeCoordinator: self.review_wiring = review_wiring self.setup_called = False self.start_kwargs: dict[str, Any] | None = None + self.new_task_callback: Any = None _FakeCoordinator.instances.append(self) def setup(self) -> None: self.setup_called = True + def set_new_task_callback(self, callback: Any) -> None: + self.new_task_callback = callback + def start_task(self, *, task_text: str, transport_name: str) -> str: self.start_kwargs = {"task_text": task_text, "transport_name": transport_name} return "thread-minted-42" diff --git a/agent-team/tests/test_ws_activation_wiring.py b/agent-team/tests/test_ws_activation_wiring.py new file mode 100644 index 0000000..b0d1bd7 --- /dev/null +++ b/agent-team/tests/test_ws_activation_wiring.py @@ -0,0 +1,160 @@ +"""Activation-wiring tests (integration branch): prove the WS seams that the +``serve`` path flips ON are actually wired, without a live Slack socket. + +Covers: +* run-team ``_build_context_provider`` returns the handbook loader (WS5 / D10). +* run-team ``_build_coordinator`` threads that context_provider into the planner + node factory. +* run-team ``_build_coordinator`` sets a ``/new-task`` callback that starts a + task on the SAME coordinator with transport_name="slack" (WS2). +* ``default_slack_listener_factory`` forwards ``new_task_callback`` to the + SlackListener (WS2). +* ``Coordinator`` stores ``new_task_callback`` / ``set_new_task_callback`` and + forwards it when it builds the default listener. +""" + +from __future__ import annotations + +import importlib.util +import sys +from pathlib import Path +from types import SimpleNamespace +from typing import Any +from unittest.mock import MagicMock + + +_AGENT_TEAM_DIR = Path(__file__).resolve().parents[1] +if str(_AGENT_TEAM_DIR) not in sys.path: + sys.path.insert(0, str(_AGENT_TEAM_DIR)) + + +def _load_run_team(): + """Import run-team.py (hyphenated, so loaded by path) as a module.""" + cli_path = _AGENT_TEAM_DIR / "run-team.py" + spec = importlib.util.spec_from_file_location("run_team_cli", cli_path) + assert spec is not None and spec.loader is not None + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + return module + + +def _dry_args(tmp_path: Path) -> SimpleNamespace: + return SimpleNamespace( + db=str(tmp_path / "agent_team.sqlite"), + transport="slack", + dry_run=True, + ) + + +# --------------------------------------------------------------------------- # +# WS5: context_provider = handbook loader, threaded into the plan node +# --------------------------------------------------------------------------- # + + +def test_build_context_provider_returns_handbook_loader() -> None: + cli = _load_run_team() + provider = cli._build_context_provider() + assert callable(provider) + # Zero-arg and returns a string (fail-safe: "" when no handbook dir). + result = provider() + assert isinstance(result, str) + + +def test_build_coordinator_threads_context_provider_into_plan_node( + tmp_path: Path, monkeypatch: Any +) -> None: + cli = _load_run_team() + from agent_team import coordinator as coord_mod + + captured: dict[str, Any] = {} + + def _spy_plan_factory(context_provider: Any = None): + captured["context_provider"] = context_provider + return lambda state: state + + monkeypatch.setattr(coord_mod, "default_plan_node_factory", _spy_plan_factory) + + coordinator = cli._build_coordinator(_dry_args(tmp_path)) + # The plan node factory the coordinator holds is the run-team lambda; calling + # it must invoke default_plan_node_factory WITH a non-None context_provider. + coordinator._build_plan_node() + assert "context_provider" in captured + assert callable(captured["context_provider"]) + + +# --------------------------------------------------------------------------- # +# WS2: /new-task callback wired to this coordinator's start_task +# --------------------------------------------------------------------------- # + + +def test_build_coordinator_wires_new_task_callback_to_start_task( + tmp_path: Path, monkeypatch: Any +) -> None: + cli = _load_run_team() + coordinator = cli._build_coordinator(_dry_args(tmp_path)) + + # Replace start_task so we can observe the callback routing without running + # the real graph. + calls: dict[str, Any] = {} + + def _fake_start_task(*, task_text: str, transport_name: str) -> str: + calls["task_text"] = task_text + calls["transport_name"] = transport_name + return "thread-xyz" + + monkeypatch.setattr(coordinator, "start_task", _fake_start_task) + + cb = coordinator._new_task_callback + assert cb is not None + thread_id = cb("fix the flaky test", "slack") + assert thread_id == "thread-xyz" + assert calls == {"task_text": "fix the flaky test", "transport_name": "slack"} + + +def test_set_new_task_callback_overrides() -> None: + from agent_team.coordinator import Coordinator + from agent_team.transport.base import Transport + + class _T(Transport): + def post_question(self, **kw: Any) -> str: # type: ignore[override] + return "q" + + def parse_answer(self, raw: Any): # type: ignore[override] + raise NotImplementedError + + coord = Coordinator(db_path=":memory:", transport=_T()) + assert coord._new_task_callback is None + sentinel = lambda t, s: "tid" # noqa: E731 + coord.set_new_task_callback(sentinel) + assert coord._new_task_callback is sentinel + + +# --------------------------------------------------------------------------- # +# WS2: factory + coordinator forward new_task_callback to the SlackListener +# --------------------------------------------------------------------------- # + + +def test_slack_listener_factory_forwards_new_task_callback( + tmp_path: Path, monkeypatch: Any +) -> None: + import agent_team.coordinator as coord_mod + + captured: dict[str, Any] = {} + + class _FakeListener: + def __init__(self, *args: Any, **kwargs: Any) -> None: + captured["new_task_callback"] = kwargs.get("new_task_callback") + + # Patch the lazily-imported SlackListener symbol. + import agent_team.transport.slack_listener as sl_mod + + monkeypatch.setattr(sl_mod, "SlackListener", _FakeListener) + + sentinel = lambda t, s: "tid" # noqa: E731 + coord_mod.default_slack_listener_factory( + transport=MagicMock(), + db_path=tmp_path / "x.sqlite", + enqueue_resume=lambda _x: None, + new_task_callback=sentinel, + ) + assert captured["new_task_callback"] is sentinel From 0689d1696c92031382d91b54700890db9ee57125 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 12:43:39 -0400 Subject: [PATCH 02/13] feat(ops): R720 WS-rollout deploy/update script (G) Idempotent attended update of the live coordinator to WS0-WS5: rsync repo + handbook, install fastapi/uvicorn, append AGENT_TEAM_API_TOKEN/SEA_HAVEN_HANDBOOK_DIR to secrev.env if absent, restart the daemon, smoke tests. HTTP API is an opt-in separate step; P3 dispatch stays inert. Snapshot-first + confirm before restart. --- agent-team/scripts/deploy-r720-ws-rollout.sh | 104 +++++++++++++++++++ 1 file changed, 104 insertions(+) create mode 100755 agent-team/scripts/deploy-r720-ws-rollout.sh diff --git a/agent-team/scripts/deploy-r720-ws-rollout.sh b/agent-team/scripts/deploy-r720-ws-rollout.sh new file mode 100755 index 0000000..811f972 --- /dev/null +++ b/agent-team/scripts/deploy-r720-ws-rollout.sh @@ -0,0 +1,104 @@ +#!/usr/bin/env bash +# deploy-r720-ws-rollout.sh — attended UPDATE of the live R720 agent-team +# coordinator to the WS0–WS5 rollout (PRs #43–#46 + the activation wiring). +# +# This is an UPDATE, not a first-time provision: P1 is already deployed per +# agent-team/DEPLOY-R720.md (repo rsynced to ~/orchestrator, venv at +# agent-team/.venv, systemd unit agent-team-coordinator.service running). +# Run this from the MAC, after the WS branches have merged to main. It rsyncs +# the new code, installs the new deps, appends the new secrets if absent, syncs +# the engineering handbook, restarts the coordinator, and smoke-tests. +# +# It is idempotent and FAILS LOUDLY. It changes a live box, so: +# 1) SNAPSHOT FIRST (Hyper-V checkpoint of sh-secrev on the R720 host). +# 2) It prompts before the restart. +# +# What goes LIVE after this (the safe, ungated seams): +# * WS1 in-process multi-model invokers (bind_multi_invoker, already wired) +# * WS5 handbook context_provider injected into the planner prompt +# * WS2 Slack /new-task -> start a task (AUTHZ-01 owner allowlist gated) +# The HTTP API (WS1) + the /delegate hook are OPTIONAL and started separately +# (see step 6). The P3 dispatch/build-verify path stays INERT (gated). +set -euo pipefail + +# ── Config (override via env) ──────────────────────────────────────────────── +BOX="${BOX:-adam@10.10.60.120}" +SSH_KEY="${SSH_KEY:-$HOME/.ssh/r720_seahaven}" +REPO_LOCAL="${REPO_LOCAL:-$HOME/Documents/repositories/orchestrator}" +HANDBOOK_LOCAL="${HANDBOOK_LOCAL:-$HOME/Documents/repositories/engineering-handbook}" +# Where the handbook lands on the box; must match SEA_HAVEN_HANDBOOK_DIR below. +HANDBOOK_REMOTE="${HANDBOOK_REMOTE:-/home/adam/.sea-haven/engineering-handbook}" +SSH="ssh -i ${SSH_KEY} ${BOX}" + +say() { printf '\n\033[1;36m== %s\033[0m\n' "$*"; } +confirm() { read -r -p "$1 [y/N] " a; [ "$a" = "y" ] || [ "$a" = "Y" ]; } + +say "Preflight" +[ -f "${SSH_KEY}" ] || { echo "missing SSH key ${SSH_KEY}"; exit 1; } +$SSH true || { echo "cannot reach ${BOX}"; exit 1; } +echo "SNAPSHOT REMINDER: take a Hyper-V checkpoint of sh-secrev on the R720 host now." +confirm "Snapshot taken and ready to update the LIVE coordinator?" || { echo "aborted"; exit 1; } + +say "1. rsync repo (Mac -> box; same excludes as the P1 runbook)" +rsync -av --exclude .env --exclude .venv --exclude .git --exclude '__pycache__' \ + "${REPO_LOCAL}/" "${BOX}:orchestrator/" + +say "2. rsync engineering handbook -> ${HANDBOOK_REMOTE} (WS5 context_provider source)" +if [ -d "${HANDBOOK_LOCAL}" ]; then + $SSH "mkdir -p ${HANDBOOK_REMOTE}" + rsync -av --delete --exclude .git "${HANDBOOK_LOCAL}/" "${BOX}:${HANDBOOK_REMOTE}/" +else + echo "WARN: ${HANDBOOK_LOCAL} not found; context_provider will return '' (fail-safe). Skipping." +fi + +say "3. Install new venv deps (fastapi/uvicorn for the WS1 HTTP API)" +# requirements.txt now pins fastapi==0.136.1 / uvicorn==0.46.0. The coordinator +# itself does not need them, but the optional HTTP API (step 6) does. +$SSH 'cd ~/orchestrator/agent-team && . .venv/bin/activate && pip install --upgrade "fastapi==0.136.1" "uvicorn==0.46.0"' + +say "4. Append new secrets to ~/secrev.env if absent (mode 600, never committed)" +# AGENT_TEAM_API_TOKEN: required only if you run the HTTP API / /delegate hook. +# SEA_HAVEN_HANDBOOK_DIR: where load_handbook_conventions() reads from. +$SSH "bash -s" <> ~/secrev.env +if grep -q '^AGENT_TEAM_API_TOKEN=' ~/secrev.env; then + echo 'AGENT_TEAM_API_TOKEN already set; leaving as-is.' +else + echo 'AGENT_TEAM_API_TOKEN NOT set. Add it now (generated on the Mac):' + echo ' echo "AGENT_TEAM_API_TOKEN=" >> ~/secrev.env && chmod 600 ~/secrev.env' + echo '(only needed for the HTTP API / auto-delegate hook; the coordinator runs without it.)' +fi +REMOTE + +say "5. Restart the coordinator daemon" +confirm "Restart agent-team-coordinator.service now?" || { echo "skipped restart"; exit 0; } +$SSH 'sudo systemctl restart agent-team-coordinator.service && sleep 2 && systemctl is-active agent-team-coordinator.service' +$SSH 'journalctl -u agent-team-coordinator.service -n 30 --no-pager' + +say "6. (OPTIONAL) HTTP API + /delegate hook — start only if you want them" +cat <<'NOTE' +The coordinator now serves WS5 context + WS2 /new-task. The WS1 HTTP API is a +SEPARATE process (api.serve(), 127.0.0.1:8765, bearer auth). To run it: + - ensure AGENT_TEAM_API_TOKEN is set in ~/secrev.env + - run: cd ~/orchestrator/agent-team && . .venv/bin/activate && \ + python3 -c "from agent_team.api import serve; serve()" + - (for persistence, add a second systemd unit; not auto-installed here.) +Then set AGENT_TEAM_API_TOKEN + AGENT_TEAM_API_URL in the Mac Claude Code env +to enable the /delegate hook (sea-haven-claude-plugin). +NOTE + +say "7. SMOKE TESTS (manual)" +cat <<'SMOKE' + a) Coordinator up: systemctl is-active agent-team-coordinator.service -> active + b) Handbook visible: cd ~/orchestrator/agent-team && . .venv/bin/activate && \ + python3 -c "from agent_team.nodes.handbook import load_handbook_conventions as h; print(bool(h()))" -> True + c) Slack /new-task: post "/new-task add a smoke-test file" in #agent-team as an + allowlisted owner -> the bot replies with a clarifying question. + d) (if API running) auth: curl -s -o /dev/null -w '%{http_code}' \ + -H "Authorization: Bearer $AGENT_TEAM_API_TOKEN" http://127.0.0.1:8765/tasks -> 405 (GET not allowed = API up + authed) +ROLLBACK: restore the pre-update Hyper-V checkpoint (one-command revert). +SMOKE +say "Done." From 2e5972c476a1295715b916d4bc27a68b147c8c0f Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 12:56:08 -0400 Subject: [PATCH 03/13] docs(integration): document WS0-WS5 components + WS-rollout deploy --- README.md | 4 +++ agent-team/DEPLOY-R720.md | 45 ++++++++++++++++++++++++ agent-team/README.md | 31 ++++++++++++++++- docs/provisioning/OPERATOR-RUNBOOK.md | 49 +++++++++++++++++++++++++++ 4 files changed, 128 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index 3f282c3..39d9b23 100644 --- a/README.md +++ b/README.md @@ -123,6 +123,10 @@ The `security-review/` subsystem is a high-recall, anti-complacency security gat See `security-review/README.md` for full detail and `security-review/DEPLOY-R720.md` for the VM runbook. +## agent-team (R720 durable SDLC pipeline) + +The `agent-team/` subsystem is a separate, durable, human-gated SDLC pipeline (LangGraph + SQLite ledger) that runs as an always-on coordinator daemon on the same `sh-secrev` R720 VM. It is distinct from this stateless router: it persists tasks across restarts and runs INTAKE → CLARIFY → PLAN → REVIEW (build/verify is deploy-gated and inert). It reuses this orchestrator's `models.py` for its non-Claude invokers, and exposes an opt-in FastAPI HTTP API (loopback, bearer auth) plus a `/delegate` Claude Code plugin hook (`sea-haven-claude-plugin/`). See `agent-team/README.md` and `agent-team/DEPLOY-R720.md`. + ## Setup 1. Install dependencies: `pip install -r requirements.txt` diff --git a/agent-team/DEPLOY-R720.md b/agent-team/DEPLOY-R720.md index 18cc9a7..71f9a8f 100644 --- a/agent-team/DEPLOY-R720.md +++ b/agent-team/DEPLOY-R720.md @@ -118,6 +118,51 @@ The unit runs `python3 run-team.py serve` from secrets from `EnvironmentFile=/home/adam/secrev.env`. `Restart=on-failure` keeps it up across transient faults; `journalctl -u` is the live log. +## 4b. WS0–WS5 rollout — UPDATE an already-deployed box + +The steps above (§1–4) are the **first-time** P1 provision. To bring an +already-deployed coordinator up to the WS0–WS5 rollout, use the attended update +script rather than re-running the manual steps: + +``` +# From the Mac, after the WS branches have merged to main, snapshot first: +agent-team/scripts/deploy-r720-ws-rollout.sh +``` + +It is an UPDATE (not a provision): it snapshots-reminds, rsyncs the new code, +rsyncs the engineering handbook to the box, installs the new deps, appends the +new secrets if absent, restarts the coordinator, and smoke-tests. It is +idempotent and fails loudly. What goes **live** after it: WS1 in-process +multi-model invokers (`bind_multi_invoker`, already wired), the WS5 handbook +`context_provider` injected into the planner prompt, and the WS2 Slack +`/new-task` command (AUTHZ-01 owner-allowlist gated). The P3 dispatch/build-verify +path stays **inert** (gated behind the `agent-apply` GitHub Environment approval). + +**New venv deps** (the coordinator does not need them; only the optional HTTP +API does) — now pinned in the root `requirements.txt`: + +| pip dep | Why | +|---|---| +| `fastapi==0.136.1` | the WS1 HTTP API app (`agent_team/api.py`) | +| `uvicorn==0.46.0` | ASGI server for `api.serve()` | + +**New env vars** — append to `~/secrev.env` (mode 600, never committed): + +- `SEA_HAVEN_HANDBOOK_DIR` — where `load_handbook_conventions()` reads the + engineering handbook (the script syncs it to `/home/adam/.sea-haven/engineering-handbook` + by default; this var must match). Fail-safe: if the dir is missing the + `context_provider` returns `""` and the planner runs without handbook context. +- `AGENT_TEAM_API_TOKEN` — bearer token for the HTTP API / `/delegate` hook + **only**. Not needed by the coordinator daemon itself. The HTTP API refuses to + start if this is unset/empty. + +**The HTTP API is a separate, opt-in process** — it is **not** started by the +coordinator daemon. Run it explicitly (`api.serve()`, binds `127.0.0.1:8765`, +bearer auth) only if you want the `/delegate` Claude Code hook or the +`POST /tasks` / `GET /tasks/{thread_id}` / `POST /orchestrator/invoke` endpoints. +The `/docs` + `/openapi` routes are disabled and it binds loopback by design (do +not change to `0.0.0.0`). See the deploy script's step 6 for how to start it. + ## 5. P1 live exit-criteria demo (§3.3.1) Demonstrate all four once the service is live. Map each to the operator commands diff --git a/agent-team/README.md b/agent-team/README.md index bf43434..8219fe1 100644 --- a/agent-team/README.md +++ b/agent-team/README.md @@ -39,6 +39,10 @@ agent-team/ graph.py # LangGraph wiring: P1 (intake→clarify→plan) + opt-in P2 # review loop + opt-in P3 build/verify subgraph invoker.py # §3.1 real Claude path (subscription-OAuth / API / Bedrock) + invoker_multi.py # WS1 in-process non-Claude invokers (GPT-4.1 / DeepSeek / + # Gemini via the orchestrator's models.py); bind_multi_invoker() + api.py # WS1 FastAPI HTTP API (bearer auth, 127.0.0.1:8765) — SEPARATE + # opt-in process (api.serve()), NOT started by the coordinator billing.py # §3.1 claude_invoke billing-mode seam ci_gate.py # §3.3.2 pure-code authenticated-Checks PASS/FAIL gate task_model.py / state_store.py @@ -52,17 +56,26 @@ agent-team/ builders.py + builders_llm.py # candidate diff (DeepSeek) — INERT, proposes only verifier.py + verifier_llm.py # ci_gate sole PASS authority; LLM = fix-proposer build_verify_subgraph.py # P3 BUILD→VERIFY topology (opt-in) + handbook.py # WS5 load_handbook_conventions (handbook seam, + # fail-safe → "" if dir missing); planner context + dispatch_invoker.py # WS3 auto-dispatch node — INERT (NOT wired live) transport/ # one adapter contract + a live impl per channel base.py # Transport ABC + QuestionSet / NormalizedAnswer - slack_adapter.py + slack_live.py + slack_listener.py # Block Kit + Socket Mode + slack_adapter.py + slack_live.py + slack_listener.py # Block Kit + Socket Mode + /new-task github_adapter.py + github_live.py + github_intake.py # issue-comment + issue intake claude_code_adapter.py + claude_code_live.py # file-drop responder + scripts/ # deploy-r720-ws-rollout.sh — attended WS0–WS5 UPDATE of the box ci/ # §3.3.2 split-job CI apply/verify workflow (DEPLOY-GATED) systemd/ # agent-team-coordinator.service (not installed) DEPLOY-R720.md # provisioning runbook (snapshot-first, rsync, tokens, demo) tests/ # pytest, one module per source module + sim harness ``` +The Claude Code plugin lives in a sibling top-level dir, `../sea-haven-claude-plugin/` +(CLAUDE.md, settings.template.json, `hooks/user_prompt_submit.py`): a +`UserPromptSubmit` hook that forwards `/delegate ` prompts from Claude Code +to the HTTP API's `POST /tasks` (env `AGENT_TEAM_API_URL` / `AGENT_TEAM_API_TOKEN`). + The top directory is kebab-case (`agent-team/`); the importable package is snake_case (`agent_team/`), per the engineering handbook. @@ -85,6 +98,22 @@ snake_case (`agent_team/`), per the engineering handbook. GPT-4.1 (review) and DeepSeek (builders) route through the local orchestrator `run.py`. Switching Claude billing is a config flip. +## WS0–WS5 rollout glossary + +The "WS-rollout" (workstreams 0–5) layered HTTP/integration surfaces onto the +P1–P4 pipeline. What is **live** vs **inert** after the rollout: + +| WS | What it adds | Live? | +|---|---|---| +| WS1 | `invoker_multi.py` (in-process GPT-4.1 / DeepSeek / Gemini via the orchestrator's `models.py`) + `api.py` (FastAPI HTTP API, bearer auth via `AGENT_TEAM_API_TOKEN`, binds `127.0.0.1:8765`, `/docs`+`/openapi` disabled, concurrency-capped) | `bind_multi_invoker()` wired in `run-team.py` `_cmd_serve` (LIVE); the **HTTP API is a separate opt-in process** (`api.serve()`), NOT started by the coordinator | +| WS5 | `nodes/handbook.py` `load_handbook_conventions` (reads `SEA_HAVEN_HANDBOOK_DIR` or `~/.sea-haven/engineering-handbook`, fail-safe → `""`); `retriever.py` `save_memory` writes to a `_box-drafts/` review queue | LIVE — the planner prompt receives the handbook via the `context_provider` seam in `run-team.py` `_build_coordinator` | +| WS2/WS0/WS4 | Slack `/new-task` slash command (AUTHZ-01 owner-allowlist gated) → `Coordinator.set_new_task_callback`; the `sea-haven-claude-plugin/` (CLAUDE.md, settings, `/delegate` `UserPromptSubmit` hook) | LIVE (`/new-task` wired in `serve`); the plugin/HTTP-API path is opt-in | +| WS3 | `nodes/dispatch_invoker.py` (auto-dispatch LangGraph node) + graph/coordinator wiring | **INERT — NOT wired live.** The `agent-apply` GitHub Environment human-approval gate is KEPT; the P3 dispatch/build-verify path stays inert pending per-task `run_id` plumbing + a CI-boundary security re-review | + +The HTTP API endpoints: `POST /tasks` (start a task), `GET /tasks/{thread_id}` +(status), `POST /orchestrator/invoke` (one-shot model invoke). See +`DEPLOY-R720.md` for the WS-rollout deploy (`scripts/deploy-r720-ws-rollout.sh`). + ## Running the tests ``` diff --git a/docs/provisioning/OPERATOR-RUNBOOK.md b/docs/provisioning/OPERATOR-RUNBOOK.md index 806b583..98f1f38 100644 --- a/docs/provisioning/OPERATOR-RUNBOOK.md +++ b/docs/provisioning/OPERATOR-RUNBOOK.md @@ -237,6 +237,55 @@ journalctl -u agent-team-coordinator.service -e | grep -i "inbound Slack listene --- +## Incident 5b — WS0–WS5 surfaces (HTTP API, /new-task, handbook context) + +These were added by the WS-rollout (see `agent-team/DEPLOY-R720.md` §4b). They +layer onto the coordinator; none of them should take down the maintenance loop. + +- **HTTP API down / unreachable** — the WS1 FastAPI app (`agent_team/api.py`) is + a **separate, opt-in process** (`api.serve()`, `127.0.0.1:8765`, bearer auth), + **not** started by the coordinator daemon. If `/delegate` from Claude Code or + `POST /tasks` over HTTP stops working, the coordinator itself is unaffected — + check the API process separately: + ```bash + curl -sS -o /dev/null -w '%{http_code}\n' \ + -H "Authorization: Bearer $AGENT_TEAM_API_TOKEN" http://127.0.0.1:8765/tasks + # 405 = API up + authed (GET not allowed on /tasks); 000 = process down; + # 401 = AGENT_TEAM_API_TOKEN mismatch (client vs ~/secrev.env). + ``` + The API refuses to start if `AGENT_TEAM_API_TOKEN` is unset/empty (logs a + `RuntimeError`). Fix the token, restart the API process. Tasks already in the + ledger are unaffected — the API is only an *intake/invoke* front door; answer + via Slack or the CLI as usual. + +- **`/new-task` Slack command not responding** — the WS2 slash command is + AUTHZ-01 owner-allowlist gated and routes through the same Socket Mode listener + as answers. If it silently does nothing, it is almost always the owner + allowlist (same failure mode as Incident 3's live-Slack path): + ```bash + journalctl -u agent-team-coordinator.service -e | grep -iE "new-task|owner|unauthorized" + # unauthorized sender / unconfigured allowlist -> fix AGENT_TEAM_SLACK_OWNER_IDS + ``` + Fallback: start the task from the CLI (`run-team.py start --task "..."`) or the + HTTP API. If the listener itself is down, see Incident 5 (inbound listener). + +- **Handbook dir missing → planner runs without handbook context** — the WS5 + `context_provider` (`load_handbook_conventions`) is **fail-safe**: if + `SEA_HAVEN_HANDBOOK_DIR` (or `~/.sea-haven/engineering-handbook`) is missing or + unreadable it returns `""` and the planner runs normally, just without handbook + conventions injected. This is **degraded, not broken** — no park, no alarm. + Confirm and restore: + ```bash + grep '^SEA_HAVEN_HANDBOOK_DIR=' ~/secrev.env + ls "$(grep '^SEA_HAVEN_HANDBOOK_DIR=' ~/secrev.env | cut -d= -f2)" # dir present + populated? + ``` + Re-sync the handbook (the deploy script does this) and restart the daemon so + the planner picks it back up. The WS3 dispatch node is **inert** (gated) and + should never appear in pipeline activity; if it does, treat as an unexpected + state and escalate. + +--- + ## Incident 6 — COMPLACENCY / COVERAGE alarms (Plane-1 checkers) These come from the nightly checker run, not the coordinator daemon (design §6.4, From 4a48f9459ccf6f36efd585ab22f7d25813e88c68 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 13:20:08 -0400 Subject: [PATCH 04/13] fix(ws2): register /new-task slash command + commands scope in Slack manifest The slack_listener handles {type:slash_commands, command:/new-task} but the app manifest declared no slash commands and no 'commands' scope, so Slack never offered /new-task (the command can't be invoked). Add the slash command + commands bot scope. Socket Mode delivers it over the socket (no request URL). APPLY: update the app A0BCC7TTU66 from this manifest + reinstall to pick up the new scope. --- agent-team/slack/agent-team-manifest.json | 13 +++++++++++-- 1 file changed, 11 insertions(+), 2 deletions(-) diff --git a/agent-team/slack/agent-team-manifest.json b/agent-team/slack/agent-team-manifest.json index 64c1e30..6345972 100644 --- a/agent-team/slack/agent-team-manifest.json +++ b/agent-team/slack/agent-team-manifest.json @@ -8,7 +8,15 @@ "bot_user": { "display_name": "agent-team", "always_online": true - } + }, + "slash_commands": [ + { + "command": "/new-task", + "description": "Start a new agent-team task (clarify -> plan -> review)", + "usage_hint": "", + "should_escape": false + } + ] }, "oauth_config": { "scopes": { @@ -17,7 +25,8 @@ "channels:history", "groups:history", "im:history", - "app_mentions:read" + "app_mentions:read", + "commands" ] } }, From 75f1dc6f753547b4bd19e46486662f64430b8480 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 13:30:13 -0400 Subject: [PATCH 05/13] fix(ws1/ws5): bootstrap orchestrator root in in-process invokers (G1); thread handbook into clarifier (G4) Gap-audit findings: - G1 (blocks-feature): make_cross_reviewer_invoker / make_fast_coder_invoker did 'from models import' without putting the orchestrator root on sys.path. The run-team serve daemon only bootstraps agent-team/, so on the live box every GPT-4.1 plan review hit ModuleNotFoundError -> review_plan's blanket except silently fail-closed to REQUEST_CHANGES (GPT-4.1 never actually ran). Both in-process invokers now call invoker_multi._ensure_orchestrator_on_path() before the deferred import. WS1 introduced this when it swapped the review default from the subprocess invoker to in-process. - G4 (degrades): the handbook context_provider was wired into the planner only; default_clarify_node_factory now accepts + forwards it, and run-team wires it into build_clarify_node too, so clarifying questions are handbook-aware. Tests: +2 regression tests (path-bootstrap, clarifier threading); _FakeCoordinator gains build_clarify_node. 1142 passed, ruff clean. --- agent-team/agent_team/coordinator.py | 12 +++- agent-team/agent_team/nodes/builders_llm.py | 7 +++ .../agent_team/nodes/review_loop_llm.py | 8 +++ agent-team/run-team.py | 15 +++-- agent-team/tests/test_run_team.py | 2 + agent-team/tests/test_ws_activation_wiring.py | 60 +++++++++++++++++++ 6 files changed, 97 insertions(+), 7 deletions(-) diff --git a/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index a8336c6..75cd7d9 100644 --- a/agent-team/agent_team/coordinator.py +++ b/agent-team/agent_team/coordinator.py @@ -187,7 +187,9 @@ def default_slack_listener_factory( ) -def default_clarify_node_factory() -> Callable[[PipelineState], PipelineState]: +def default_clarify_node_factory( + context_provider: "Callable[[], str] | None" = None, +) -> Callable[[PipelineState], PipelineState]: """Build the live Claude-backed clarifier node (§3.3, §7.1 P1). Composes the two committed leaves: the Claude clarifier callables @@ -206,7 +208,13 @@ def default_clarify_node_factory() -> Callable[[PipelineState], PipelineState]: from agent_team.nodes.clarifier import make_clarifier_node from agent_team.nodes.clarifier_llm import build_claude_clarifier_callables - assess_confidence, generate_questions = build_claude_clarifier_callables() + # WS5: thread the handbook/memory context_provider into the clarifier too + # (not just the planner) so clarifying questions are handbook-aware. Only + # forwarded when non-None to preserve the byte-identical default behavior. + kwargs: dict[str, Any] = {} + if context_provider is not None: + kwargs["context_provider"] = context_provider + assess_confidence, generate_questions = build_claude_clarifier_callables(**kwargs) return make_clarifier_node( assess_confidence=assess_confidence, generate_questions=generate_questions, diff --git a/agent-team/agent_team/nodes/builders_llm.py b/agent-team/agent_team/nodes/builders_llm.py index 7f3ad78..812ddfd 100644 --- a/agent-team/agent_team/nodes/builders_llm.py +++ b/agent-team/agent_team/nodes/builders_llm.py @@ -167,6 +167,13 @@ def make_fast_coder_invoker() -> "BuildCallable": """ def _invoke(instruction: str) -> str: + # Bootstrap the orchestrator root onto sys.path (models.py lives there; + # the run-team serve daemon does not add it) before the deferred import, + # mirroring make_cross_reviewer_invoker. Without it `from models import` + # raises ModuleNotFoundError in the daemon. + from agent_team.invoker_multi import _ensure_orchestrator_on_path # noqa: PLC0415 + + _ensure_orchestrator_on_path() from models import get_fast_coder # noqa: PLC0415 - intentional deferred import coder = get_fast_coder() diff --git a/agent-team/agent_team/nodes/review_loop_llm.py b/agent-team/agent_team/nodes/review_loop_llm.py index 3d8de87..50c4d62 100644 --- a/agent-team/agent_team/nodes/review_loop_llm.py +++ b/agent-team/agent_team/nodes/review_loop_llm.py @@ -210,6 +210,14 @@ def make_cross_reviewer_invoker() -> PlanReviewer: """ def _invoke(prompt: str, *, config: Any = None, **_kw: Any) -> str: + # The orchestrator root (where models.py lives) is NOT on sys.path in + # the run-team serve daemon (it only bootstraps agent-team/). Without + # this, `from models import` raises ModuleNotFoundError → review_plan's + # blanket except silently fails-closed to REQUEST_CHANGES, so GPT-4.1 + # never actually runs. Bootstrap the root before the deferred import. + from agent_team.invoker_multi import _ensure_orchestrator_on_path # noqa: PLC0415 + + _ensure_orchestrator_on_path() # Deferred import: keep the orchestrator package out of module import. from models import get_cross_reviewer # noqa: PLC0415 diff --git a/agent-team/run-team.py b/agent-team/run-team.py index 33d914f..9248ccd 100644 --- a/agent-team/run-team.py +++ b/agent-team/run-team.py @@ -529,16 +529,18 @@ def _build_coordinator(args: argparse.Namespace) -> Any: """ from agent_team.coordinator import ( Coordinator, + default_clarify_node_factory, default_plan_node_factory, default_review_wiring, ) transport = _build_transport(args) - # WS5 (D10): inject the Sea Haven handbook conventions into the planner prompt - # via the context_provider seam. The provider is the zero-arg handbook loader, - # which is itself fail-safe (returns "" when the handbook dir is absent), and - # planner.plan_node additionally swallows provider errors — so this never - # affects a run where the handbook is unavailable. + # WS5 (D10): inject the Sea Haven handbook conventions into BOTH the clarifier + # and the planner prompts via the context_provider seam. The provider is the + # zero-arg handbook loader, which is itself fail-safe (returns "" when the + # handbook dir is absent), and the clarifier/planner additionally swallow + # provider errors — so this never affects a run where the handbook is + # unavailable. context_provider = _build_context_provider() # Production runs the full P2 graph: the wrapped real planner + the bound @@ -547,6 +549,9 @@ def _build_coordinator(args: argparse.Namespace) -> Any: coordinator = Coordinator( db_path=args.db, transport=transport, + build_clarify_node=lambda: default_clarify_node_factory( + context_provider=context_provider + ), build_plan_node=lambda: default_plan_node_factory( context_provider=context_provider ), diff --git a/agent-team/tests/test_run_team.py b/agent-team/tests/test_run_team.py index 7ca0ce7..9d180ae 100644 --- a/agent-team/tests/test_run_team.py +++ b/agent-team/tests/test_run_team.py @@ -644,6 +644,7 @@ class _FakeCoordinator: *, db_path: Any, transport: Any, + build_clarify_node: Any = None, build_plan_node: Any = None, review_wiring: Any = None, ) -> None: @@ -651,6 +652,7 @@ class _FakeCoordinator: self.transport = transport # 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 self.build_plan_node = build_plan_node self.review_wiring = review_wiring self.setup_called = False diff --git a/agent-team/tests/test_ws_activation_wiring.py b/agent-team/tests/test_ws_activation_wiring.py index b0d1bd7..9988d9c 100644 --- a/agent-team/tests/test_ws_activation_wiring.py +++ b/agent-team/tests/test_ws_activation_wiring.py @@ -158,3 +158,63 @@ def test_slack_listener_factory_forwards_new_task_callback( new_task_callback=sentinel, ) assert captured["new_task_callback"] is sentinel + + +# --------------------------------------------------------------------------- # +# WS5 (G4): context_provider is threaded into the CLARIFIER too, not just plan +# --------------------------------------------------------------------------- # + + +def test_build_coordinator_threads_context_provider_into_clarify_node( + tmp_path: Path, monkeypatch: Any +) -> None: + cli = _load_run_team() + from agent_team import coordinator as coord_mod + + captured: dict[str, Any] = {} + + def _spy_clarify_factory(context_provider: Any = None): + captured["context_provider"] = context_provider + return lambda state: state + + monkeypatch.setattr(coord_mod, "default_clarify_node_factory", _spy_clarify_factory) + + coordinator = cli._build_coordinator(_dry_args(tmp_path)) + coordinator._build_clarify_node() + assert "context_provider" in captured + assert callable(captured["context_provider"]) + + +# --------------------------------------------------------------------------- # +# WS1 (G1): the in-process cross-reviewer bootstraps the orchestrator root onto +# sys.path before importing `models` (else the daemon silently REQUEST_CHANGES). +# --------------------------------------------------------------------------- # + + +def test_cross_reviewer_invoker_bootstraps_orchestrator_path(monkeypatch: Any) -> None: + import types + + from agent_team.invoker_multi import _ensure_orchestrator_on_path # noqa: F401 + from agent_team.nodes.review_loop_llm import make_cross_reviewer_invoker + + # The orchestrator root is parents[2] of invoker_multi.py. + import agent_team.invoker_multi as im + + root = str(Path(im.__file__).resolve().parents[2]) + + # Simulate the daemon: root NOT on sys.path. Inject a fake `models` so the + # deferred import resolves without real provider keys — the point is to + # prove the bootstrap runs (root re-added) BEFORE the import. + monkeypatch.setattr(sys, "path", [p for p in sys.path if p != root]) + fake_models = types.ModuleType("models") + fake_reviewer = MagicMock() + fake_reviewer.invoke.return_value = MagicMock(content="APPROVE") + fake_models.get_cross_reviewer = lambda: fake_reviewer # type: ignore[attr-defined] + monkeypatch.setitem(sys.modules, "models", fake_models) + + invoker = make_cross_reviewer_invoker() + out = invoker("review this plan") + assert out == "APPROVE" + assert root in sys.path, ( + "invoker must bootstrap the orchestrator root onto sys.path" + ) From 478bd7f90b4c857e98f4609e131c4169822baddc Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 13:37:11 -0400 Subject: [PATCH 06/13] fix(ops): deploy installs full requirements.txt (incl non-Claude model stack) The agent-team venv was missing langchain-anthropic/-openai/-google-genai/ -community, so models.py failed to import and the in-process GPT-4.1 review / Gemini scan / DeepSeek build silently fail-closed to REQUEST_CHANGES (the non-Claude models never ran on the box). Step 3 now installs the full pinned requirements.txt into the venv instead of just fastapi/uvicorn. Installed + verified live on the box: all three model factories construct. --- agent-team/scripts/deploy-r720-ws-rollout.sh | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/agent-team/scripts/deploy-r720-ws-rollout.sh b/agent-team/scripts/deploy-r720-ws-rollout.sh index 811f972..aab77a6 100755 --- a/agent-team/scripts/deploy-r720-ws-rollout.sh +++ b/agent-team/scripts/deploy-r720-ws-rollout.sh @@ -51,10 +51,15 @@ else echo "WARN: ${HANDBOOK_LOCAL} not found; context_provider will return '' (fail-safe). Skipping." fi -say "3. Install new venv deps (fastapi/uvicorn for the WS1 HTTP API)" -# requirements.txt now pins fastapi==0.136.1 / uvicorn==0.46.0. The coordinator -# itself does not need them, but the optional HTTP API (step 6) does. -$SSH 'cd ~/orchestrator/agent-team && . .venv/bin/activate && pip install --upgrade "fastapi==0.136.1" "uvicorn==0.46.0"' +say "3. Install venv deps from the pinned requirements.txt" +# Install the FULL pinned set into the agent-team venv. This includes the +# non-Claude model stack (langchain-anthropic/-openai/-google-genai/-community) +# that the in-process invokers (WS1: GPT-4.1 review, Gemini scan, DeepSeek build) +# import via models.py — WITHOUT these, models.py fails to import and the review +# loop silently fail-closes to REQUEST_CHANGES (the non-Claude models never run). +# Also brings fastapi/uvicorn (WS1 HTTP API). Leaves the venv-only deps that are +# NOT in requirements.txt (claude-agent-sdk, slack_sdk, slack_bolt) untouched. +$SSH 'cd ~/orchestrator/agent-team && . .venv/bin/activate && pip install --upgrade -r ~/orchestrator/requirements.txt' say "4. Append new secrets to ~/secrev.env if absent (mode 600, never committed)" # AGENT_TEAM_API_TOKEN: required only if you run the HTTP API / /delegate hook. From 0945d603382f0facacc091fb856e75a74a7fa417 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 13:43:42 -0400 Subject: [PATCH 07/13] fix(ws2): register @app.command(/new-task) so Slack slash command is acked 'the app did not respond': serve() registered @app.action/@app.event but NO @app.command handler, so Bolt never acked the /new-task slash command within Slack's ~3s deadline. Add an @app.command(/new-task) handler that ack()s first, re-stamps type:slash_commands onto Bolt's inner command body (Bolt strips the Socket Mode envelope type that _is_new_task_command/_discriminating_type expect), then forwards to handle_event (AUTHZ-01 + new-task dispatch). The resulting payload shape is the one already covered by test_new_task_calls_callback_*. --- agent-team/agent_team/transport/slack_listener.py | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/agent-team/agent_team/transport/slack_listener.py b/agent-team/agent_team/transport/slack_listener.py index a35b305..7df310d 100644 --- a/agent-team/agent_team/transport/slack_listener.py +++ b/agent-team/agent_team/transport/slack_listener.py @@ -523,6 +523,21 @@ class SlackListener: def _on_mention(body: Mapping[str, Any]) -> None: _forward(body) + # Slash commands (e.g. /new-task) MUST be ack()'d within ~3s or Slack + # shows "the app did not respond". Bolt delivers the INNER command + # payload (command / text / user_id / channel_id) WITHOUT the Socket + # Mode envelope's ``type``, so re-stamp ``type: "slash_commands"`` to + # match what handle_event's _is_new_task_command + _discriminating_type + # expect, then forward (handle_event runs AUTHZ-01 + the new-task + # dispatch). Ack FIRST so the 3s deadline is met even though start_task + # then runs the clarifier to the first gate synchronously. + @app.command(_NEW_TASK_COMMAND) + def _on_command(ack: Any, body: Mapping[str, Any]) -> None: + ack() + payload = dict(body) + payload.setdefault("type", "slash_commands") + _forward(payload) + handler = SocketModeHandler(app, self._app_token) self._socket_handler = handler handler.start() From a6fd1cfc752bd9d4c79e6f158cf407454e17e7f4 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 13:53:39 -0400 Subject: [PATCH 08/13] fix(intake): seed the task description into graph state (was silently dropped) /new-task (and every intake: GitHub issue, /sh-assign-task) reached the clarifier with NO description -> the clarifier asked 'no task description provided'. Root cause: coordinator.start_task only LOGGED task_text (a P1-era decision when the deterministic clarifier didn't consume a description), graph.start_task took no task arg, and PipelineState/TaskRecord had no 'task' channel at all. Fix: add a first-class 'task' field to PipelineState + TaskRecord (+ round-trip in task_from_dict); graph.start_task seeds task into the initial invoke (persists through intake_node's partial-state return into CLARIFY); coordinator.start_task passes task=task_text. The clarifier already reads state['task'] via _task_description, so it now sees the real description. Test: start_task(task='build a login form') -> suspended CLARIFY state carries task. 1143 passed, ruff clean. --- agent-team/agent_team/coordinator.py | 24 ++++++++++++------------ agent-team/agent_team/graph.py | 9 +++++++++ agent-team/agent_team/task_model.py | 7 +++++++ agent-team/tests/test_graph.py | 11 +++++++++++ 4 files changed, 39 insertions(+), 12 deletions(-) diff --git a/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index 75cd7d9..312ff62 100644 --- a/agent-team/agent_team/coordinator.py +++ b/agent-team/agent_team/coordinator.py @@ -624,17 +624,15 @@ class Coordinator: :func:`agent_team.responder.notify_question` (ledger row OPEN first, then transport post). Returns the minted ``thread_id``. - **Intake-seed decision (P1).** ``agent_team.graph.start_task`` builds its - own INTAKE seed and accepts only ``thread_id`` / ``transport`` — it takes - no task-description argument, and a value pre-seeded onto the START - checkpoint via ``update_state`` is overwritten by its own seed invoke - (and a post-suspend ``update_state`` clears the pending interrupt, which - would break the human gate). So for P1 the ``task_text`` is intake - metadata held coordinator-side (logged) rather than written into - ``PipelineState``: the deterministic P1 clarifier does not consume a task - description anyway, and threading it into the graph state is a later phase - that extends the committed ``start_task`` seed contract. We keep it - minimal rather than reach past that contract or disturb the gate. + **Intake-seed (task description).** ``task_text`` is passed to + ``agent_team.graph.start_task`` as ``task=`` so it is written into the + seed ``PipelineState`` (the LLM clarifier reasons about it; an empty + description makes the clarifier ask for one). It is seeded directly into + the initial invoke — NOT via ``update_state`` (a pre-seed via + ``update_state`` would be overwritten by the seed invoke, and a + post-suspend ``update_state`` would clear the pending interrupt and break + the human gate). ``intake_node`` returns only a partial state + (status/phase), so the seeded ``task`` channel persists into CLARIFY. """ if self._graph is None: raise RuntimeError("Coordinator.start_task called before setup()") @@ -647,7 +645,9 @@ class Coordinator: transport_name, task_text[:200].replace("\n", "\\n").replace("\r", "\\r"), ) - thread_id, _state = graph_mod.start_task(self._graph, transport=transport_name) + thread_id, _state = graph_mod.start_task( + self._graph, transport=transport_name, task=task_text + ) question = graph_mod.pending_question(self._graph, thread_id=thread_id) if question is None: diff --git a/agent-team/agent_team/graph.py b/agent-team/agent_team/graph.py index 7c4b996..540d547 100644 --- a/agent-team/agent_team/graph.py +++ b/agent-team/agent_team/graph.py @@ -591,6 +591,7 @@ def start_task( *, thread_id: str | None = None, transport: str = "", + task: str = "", ) -> tuple[str, PipelineState]: """Start a new pipeline task and run it up to the first human gate (§3.3). @@ -600,6 +601,13 @@ def start_task( the checkpointed snapshot after the suspend (its ``__interrupt__`` carries the pending question-set, surfaced by :func:`pending_question`). + ``task`` is the intake description (e.g. the Slack ``/new-task`` text or a + GitHub issue body). It is written into the seed ``PipelineState`` so the + clarifier can reason about it; ``intake_node`` returns only a partial state + (status/phase), so the seeded ``task`` channel persists into CLARIFY. An + empty ``task`` (the default) seeds no description — the clarifier then asks + for one. + The graph MUST be compiled with a checkpointer for the suspend to persist; an uncheckpointed graph would run straight through without honouring the interrupt. @@ -610,6 +618,7 @@ def start_task( thread_id=tid, status=TaskStatus.ACTIVE.value, current_phase=_phase_value(Phase.INTAKE), + task=task, qa_history=[], transport=transport, created_at=now, diff --git a/agent-team/agent_team/task_model.py b/agent-team/agent_team/task_model.py index 2e09a91..9408ed7 100644 --- a/agent-team/agent_team/task_model.py +++ b/agent-team/agent_team/task_model.py @@ -87,6 +87,8 @@ class TaskRecord: thread_id: str status: TaskStatus current_phase: Phase + # Intake task description (mirrors PipelineState.task). + task: str = "" qa_history: list[Any] = field(default_factory=list) plan: dict[str, Any] | None = None review_verdicts: list[Any] = field(default_factory=list) @@ -108,6 +110,10 @@ class PipelineState(TypedDict, total=False): thread_id: str status: str current_phase: str + # The intake task description (Slack /new-task text, GitHub issue body, etc.). + # Seeded by graph.start_task and read by the clarifier/planner; a first-class + # channel so the seeded value persists across node transitions. + task: str qa_history: list[Any] plan: dict[str, Any] | None review_verdicts: list[Any] @@ -133,6 +139,7 @@ def task_from_dict(data: dict[str, Any]) -> TaskRecord: thread_id=data["thread_id"], status=TaskStatus(data["status"]), current_phase=Phase(data["current_phase"]), + task=data.get("task", ""), qa_history=list(data.get("qa_history", [])), plan=data.get("plan"), review_verdicts=list(data.get("review_verdicts", [])), diff --git a/agent-team/tests/test_graph.py b/agent-team/tests/test_graph.py index 2a29dbc..48d0b97 100644 --- a/agent-team/tests/test_graph.py +++ b/agent-team/tests/test_graph.py @@ -136,6 +136,17 @@ def test_start_task_suspends_on_human_gate(compiled) -> None: assert payload["deadline"] +def test_start_task_seeds_task_description_into_state(compiled) -> None: + # Regression: the intake description (Slack /new-task text, GitHub issue body) + # must reach the graph state so the clarifier can reason about it. It is + # seeded into the initial invoke and must persist through INTAKE into the + # suspended CLARIFY snapshot (intake_node returns only a partial state). + _thread_id, state = start_task( + compiled, transport="slack", task="build a login form" + ) + assert state.get("task") == "build a login form" + + def test_pending_question_carries_foundation_questionset(compiled) -> None: thread_id, _ = start_task(compiled, transport="slack") payload = pending_question(compiled, thread_id=thread_id) From f0a2dc27a57f98aad193f3d4639d5c08d5aea8d3 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 14:17:08 -0400 Subject: [PATCH 09/13] feat(notify): Slack lifecycle notifications + deliver multi-turn questions MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- agent-team/agent_team/coordinator.py | 85 +++++++++++++++++++++++- agent-team/run-team.py | 54 +++++++++++++++ agent-team/tests/test_coordinator.py | 98 ++++++++++++++++++++++++++++ agent-team/tests/test_run_team.py | 4 ++ 4 files changed, 240 insertions(+), 1 deletion(-) 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 From 655f6d80f4e6d818be8daf2aa2d831c8935da5fe Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 14:35:21 -0400 Subject: [PATCH 10/13] =?UTF-8?q?feat(notify):=20richer=20park/lifecycle?= =?UTF-8?q?=20messages=20=E2=80=94=20task=20description=20+=20phase=20+=20?= =?UTF-8?q?blocker?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Park notifications were opaque ('Task 6c3c3202 parked — needs your attention'): no idea what the task is, where it got to, or what's blocking it. Now each message names the task DESCRIPTION (not just the short id), the PHASE it reached, and the actual BLOCKER — _summarize_blocker() pulls the last review_verdict's findings (the GPT-4.1 REQUEST_CHANGES text, collapsed + truncated), falling back to 'no plan built' / 'revision cap hit'. Applies to parked + needs-more-input + plan-ready emits. +1 test (description + phase + blocker present). 1148 passed. --- agent-team/agent_team/coordinator.py | 66 +++++++++++++++++++++++----- agent-team/tests/test_coordinator.py | 37 ++++++++++++++++ 2 files changed, 92 insertions(+), 11 deletions(-) diff --git a/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index 16af0ea..061421a 100644 --- a/agent-team/agent_team/coordinator.py +++ b/agent-team/agent_team/coordinator.py @@ -779,6 +779,19 @@ class Coordinator: continue seen.add(thread_id) short = thread_id[:8] + + # Read the settled state once so every message can name WHAT the task + # is (description), not just an opaque thread id. + try: + snap = self._graph.get_state(graph_mod.thread_config(thread_id)) + values = getattr(snap, "values", {}) or {} + except Exception: # noqa: BLE001 + values = {} + desc = str(values.get("task") or "").strip() or "(no description)" + if len(desc) > 90: + desc = desc[:90] + "…" + label = f'"{desc}" (`{short}`)' + try: question = graph_mod.pending_question(self._graph, thread_id=thread_id) except Exception: # noqa: BLE001 - never let a status check break the loop @@ -804,23 +817,54 @@ class Coordinator: _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.") + self._emit( + f"❓ {label} — needs more input. A new clarifying question was " + "posted above; reply in its thread." + ) 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 + # No pending question: the task settled. Distinguish parked vs done, + # and say WHERE it got to and WHAT is blocking it. + status = values.get("status") + phase = str(values.get("current_phase") or "unknown") if status == TaskStatus.PARKED.value: + blocker = self._summarize_blocker(values) self._emit( - f"⚠️ Task {short}: parked — needs your attention " - "(plan/review escalation). Re-assign or steer it to resume." + f"⚠️ PARKED — {label}\n" + f"• Reached phase: {phase}\n" + f"• What's blocking it: {blocker}\n" + "• Needs your review — the plan could not be auto-approved. " + "Re-assign with more detail, or adjust the requirement to unblock." ) else: - self._emit(f"✅ Task {short}: plan ready for review.") + self._emit(f"✅ {label} — plan ready for review (phase: {phase}).") + + @staticmethod + def _summarize_blocker(values: "dict[str, Any]") -> str: + """Human-readable reason a task parked, from the last review verdict. + + Pulls the most recent ``review_verdicts`` entry's findings (the GPT-4.1 + REQUEST_CHANGES text) — collapsed + truncated for a Slack line — so the + human sees WHAT stopped it, not just "escalation". Falls back to a generic + reason when there is no verdict (e.g. an unparseable-plan park before + review ever ran). + """ + verdicts = values.get("review_verdicts") or [] + if verdicts: + last = verdicts[-1] + if isinstance(last, dict): + findings = str( + last.get("findings") or last.get("verdict") or "" + ).strip() + if findings: + summary = " ".join(findings.split()) + return summary[:300] + ("…" if len(summary) > 300 else "") + if not values.get("plan"): + return "the planner could not produce a usable plan (no plan was built)." + return ( + "the plan→review loop hit its revision cap (the reviewer kept " + "requesting changes without converging)." + ) def tick(self) -> list[ResumeResult]: """One maintenance pass: deadline sweep + park policy, then drain (§3.3.1). diff --git a/agent-team/tests/test_coordinator.py b/agent-team/tests/test_coordinator.py index b305552..ccb9e03 100644 --- a/agent-team/tests/test_coordinator.py +++ b/agent-team/tests/test_coordinator.py @@ -1010,3 +1010,40 @@ def test_emit_swallows_notify_failure(db_path: Path) -> None: coord = _make_coordinator(db_path) coord._notify = _boom coord._emit("anything") # must not raise + + +def test_parked_message_includes_description_and_blocker( + 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, + "current_phase": "review", + "task": "add a smoke-test file in agent-team/tests", + "review_verdicts": [ + { + "verdict": "request_changes", + "findings": "Missing a final review/commit phase; lint runs before tests.", + } + ], + } + + monkeypatch.setattr(coord._graph, "get_state", lambda _cfg: _Snap()) + coord._post_resume_followups([_resume_result("c0ffee01abcd")]) + assert len(msgs) == 1 + m = msgs[0] + assert "add a smoke-test file" in m # WHAT the task is + assert "review" in m # WHERE it got to + assert "Missing a final review/commit phase" in m # WHY it's blocked From e7ef4387e66822e562d95ae97e709f94849796e1 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 14:53:21 -0400 Subject: [PATCH 11/13] fix(pipeline): single-shot Claude calls + planner actually reads review findings MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two root causes behind 'every task parks, and slowly': 1. SINGLE-SHOT INVOKER: subscription Claude calls ran as 40-turn, tool-enabled agentic sessions (--max-turns 40, $2 budget) for what are pure reasoning->JSON completions — minutes-long, and Claude wandered/returned unparseable output. _DEFAULT_MAX_TURNS 40->1 + allowed_tools=[] -> fast deterministic single turn. 2. PLANNER FEEDBACK KEY MISMATCH: the review stage writes verdict/findings, but _format_review_feedback read decision/notes/comment (never present) -> the planner re-planned with EMPTY feedback, re-introduced the rejected flaw ('assumptions persist'), hit the review cap, parked. Now reads verdict/findings (old keys kept as fallback) so GPT-4.1's objections reach the re-plan. Plus: park notifications infer the phase the task was IN (review/plan/clarify) instead of the terminal 'parked'. Tests: planner real-verdict-keys regression + coordinator phase-inference. 1149 pass. --- agent-team/agent_team/coordinator.py | 10 +++++++++ agent-team/agent_team/invoker.py | 12 ++++++++++- agent-team/agent_team/nodes/planner.py | 17 +++++++++++++-- agent-team/tests/test_coordinator.py | 30 ++++++++++++++++++++++++++ agent-team/tests/test_planner.py | 25 +++++++++++++++++++-- 5 files changed, 89 insertions(+), 5 deletions(-) diff --git a/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index 061421a..3500878 100644 --- a/agent-team/agent_team/coordinator.py +++ b/agent-team/agent_team/coordinator.py @@ -827,6 +827,16 @@ class Coordinator: # and say WHERE it got to and WHAT is blocking it. status = values.get("status") phase = str(values.get("current_phase") or "unknown") + # When a task parks, current_phase is the terminal "parked" — not + # useful. Infer the phase it was IN when it escalated, so the human + # sees WHERE it died (review > plan > clarify by what state exists). + if phase == "parked": + if values.get("review_verdicts"): + phase = "review" + elif values.get("plan"): + phase = "plan" + else: + phase = "clarify" if status == TaskStatus.PARKED.value: blocker = self._summarize_blocker(values) self._emit( diff --git a/agent-team/agent_team/invoker.py b/agent-team/agent_team/invoker.py index 58d11b3..5778173 100644 --- a/agent-team/agent_team/invoker.py +++ b/agent-team/agent_team/invoker.py @@ -49,7 +49,13 @@ API_MODEL = "claude-sonnet-4-6" # Default per-call agent budget for the headless subscription path, in USD. _DEFAULT_BUDGET_USD = 2.0 -_DEFAULT_MAX_TURNS = 40 +# The agent-team's subscription Claude calls (clarify confidence/questions, +# planner) are SINGLE-SHOT reasoning→JSON completions, NOT agentic sessions. A +# 40-turn, tool-enabled session made each call take minutes and let Claude wander +# (use tools / explore) and return output the planner/clarifier couldn't parse → +# spurious parks. One turn + no tools = a fast, deterministic completion. A +# genuinely agentic caller (e.g. a future fixer) overrides max_turns/allowed_tools. +_DEFAULT_MAX_TURNS = 1 # --------------------------------------------------------------------------- # @@ -104,6 +110,10 @@ async def _collect_subscription_text( model=model, max_turns=max_turns, max_budget_usd=budget_usd, + # No tools: these are pure reasoning→JSON completions. Disallowing tools + # keeps the call a single deterministic turn (no repo exploration / tool + # loops that produce slow, unparseable output). + allowed_tools=[], ) texts: list[str] = [] diff --git a/agent-team/agent_team/nodes/planner.py b/agent-team/agent_team/nodes/planner.py index db11cbf..6d55b65 100644 --- a/agent-team/agent_team/nodes/planner.py +++ b/agent-team/agent_team/nodes/planner.py @@ -122,8 +122,21 @@ def _format_review_feedback(review_verdicts: list[Any]) -> str: lines: list[str] = [] for idx, verdict in enumerate(review_verdicts, start=1): if isinstance(verdict, dict): - decision = str(verdict.get("decision", "")).strip() - notes = str(verdict.get("notes") or verdict.get("comment") or "").strip() + # The review stage (review_loop.ReviewResult.to_dict) writes the keys + # "verdict" (decision) and "findings" (the reviewer's objections). + # Read those FIRST — the old "decision"/"notes"/"comment" keys never + # existed on a real verdict, so the planner re-planned with EMPTY + # feedback and kept re-introducing the rejected flaw ("assumptions + # persist" → cap → park). Old keys kept as fallbacks for safety. + decision = str( + verdict.get("verdict") or verdict.get("decision") or "" + ).strip() + notes = str( + verdict.get("findings") + or verdict.get("notes") + or verdict.get("comment") + or "" + ).strip() lines.append(f"Review {idx} [{decision}]: {notes}".rstrip()) else: lines.append(f"Review {idx}: {str(verdict).strip()}") diff --git a/agent-team/tests/test_coordinator.py b/agent-team/tests/test_coordinator.py index ccb9e03..66a6d2b 100644 --- a/agent-team/tests/test_coordinator.py +++ b/agent-team/tests/test_coordinator.py @@ -1047,3 +1047,33 @@ def test_parked_message_includes_description_and_blocker( assert "add a smoke-test file" in m # WHAT the task is assert "review" in m # WHERE it got to assert "Missing a final review/commit phase" in m # WHY it's blocked + + +def test_parked_message_infers_phase_when_current_phase_is_parked( + db_path: Path, monkeypatch: Any +) -> None: + # current_phase is the terminal "parked"; the message should report the phase + # the task was IN (review, since verdicts exist), not "parked". + 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, + "current_phase": "parked", + "task": "do a thing", + "review_verdicts": [{"verdict": "request_changes", "findings": "nope"}], + } + + monkeypatch.setattr(coord._graph, "get_state", lambda _cfg: _Snap()) + coord._post_resume_followups([_resume_result("abcd1234ef00")]) + assert "Reached phase: review" in msgs[0] + assert "Reached phase: parked" not in msgs[0] diff --git a/agent-team/tests/test_planner.py b/agent-team/tests/test_planner.py index 4b6c39a..988ddad 100644 --- a/agent-team/tests/test_planner.py +++ b/agent-team/tests/test_planner.py @@ -172,19 +172,40 @@ def test_build_prompt_no_task_uses_placeholder() -> None: def test_build_prompt_loopback_includes_feedback_and_prior_plan() -> None: + # Use the REAL verdict shape that review_loop.ReviewResult.to_dict() writes + # ("verdict" + "findings") — NOT the old "decision"/"notes" keys, which never + # existed on a real verdict and silently produced empty re-plan feedback. state = _state( plan={"task": "t", "phases": [{"name": "old", "steps": ["x"]}]}, review_verdicts=[ - {"decision": "REQUEST_CHANGES", "notes": "Phase 1 missing rollback."} + { + "verdict": "request_changes", + "outcome": "loop_back", + "round_index": 1, + "findings": "Phase 1 missing rollback.", + } ], ) prompt = build_plan_prompt(state) assert "Reviewer feedback" in prompt - assert "Phase 1 missing rollback." in prompt + assert "Phase 1 missing rollback." in prompt # the findings reached the re-plan assert "Previous plan" in prompt assert '"old"' in prompt +def test_review_feedback_reads_real_verdict_keys() -> None: + # Regression: the planner must read the producer's keys (verdict/findings). + # The bug read decision/notes/comment -> empty feedback -> the planner kept + # re-introducing the rejected flaw -> review cap -> park. + from agent_team.nodes.planner import _format_review_feedback + + rendered = _format_review_feedback( + [{"verdict": "request_changes", "findings": "Add a teardown fixture."}] + ) + assert "Add a teardown fixture." in rendered + assert "request_changes" in rendered + + def test_build_prompt_no_feedback_omits_review_sections() -> None: prompt = build_plan_prompt(_state(plan={"task": "t"})) assert "Reviewer feedback" not in prompt From 3847e43ba32ffa8e70c7a29e48b23e770ecba138 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 15:11:19 -0400 Subject: [PATCH 12/13] =?UTF-8?q?feat(agent-team):=20one=20Slack=20thread?= =?UTF-8?q?=20per=20task=20=E2=80=94=20root=20"Task=20received"=20message?= =?UTF-8?q?=20+=20threaded=20questions/milestones?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit WS Slack-UX Feature 1. A /new-task task now maps to ONE Slack thread instead of several top-level messages. - /new-task posts an immediate root "📥 Task received: …" ack and captures its ts (root_ts); this is the instant acknowledgement. - root_ts is plumbed into start: new PipelineState/TaskRecord channel slack_thread_ts, seeded by graph.start_task and threaded through Coordinator.start_task. The NewTaskCallback is now (task_text, via, root_ts). - All clarifier questions for the task post as THREADED REPLIES under root_ts (chat.postMessage thread_ts=root_ts), and each question's ledger channel_ref is set to root_ts (NOT the reply's own ts). Because answer-mapping resolves a reply via find_open_question_by_channel_ref(thread_ts), a reply in the root thread (thread_ts==root_ts) maps to the task's currently-open question with NO change to the mapping logic or the first-answer-wins CAS. The open-only partial-unique index still holds (one open question per task at a time). - Lifecycle milestones (parked / plan-ready / needs-input) and follow-up questions thread under root_ts too; the notify sink gained an optional thread_ts kwarg (degrades to top-level on a sink that doesn't accept it). notify failures still never break tick. - SlackTransport.post_question + the live poster accept/forward thread_ts. - No root_ts (non-/new-task origin) ⇒ top-level posts exactly as before. AUTHZ-01 (owner-allowlist-first, fail-closed) and the atomic open→answered compare-and-set are unchanged. Adds plumbing for the inbound-ack reactor seam used by Feature 2 (dormant until a reactor is injected). Tests cover thread_ts forwarding, channel_ref=root_ts, graph seeding, and coordinator threading. --- agent-team/agent_team/coordinator.py | 118 ++++++++++++-- agent-team/agent_team/graph.py | 14 ++ agent-team/agent_team/responder.py | 27 +++- agent-team/agent_team/task_model.py | 10 ++ agent-team/agent_team/transport/base.py | 7 + .../agent_team/transport/slack_adapter.py | 29 ++++ .../agent_team/transport/slack_listener.py | 153 ++++++++++++++++-- agent-team/agent_team/transport/slack_live.py | 40 ++++- agent-team/run-team.py | 6 +- agent-team/tests/test_coordinator.py | 101 +++++++++++- agent-team/tests/test_graph.py | 22 +++ agent-team/tests/test_responder.py | 53 +++++- agent-team/tests/test_slack_adapter.py | 40 +++++ .../test_ws0_ws2_ws4_plugin_slack_hook.py | 18 ++- agent-team/tests/test_ws_activation_wiring.py | 19 ++- 15 files changed, 605 insertions(+), 52 deletions(-) diff --git a/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index 3500878..74ffc45 100644 --- a/agent-team/agent_team/coordinator.py +++ b/agent-team/agent_team/coordinator.py @@ -152,7 +152,7 @@ def default_slack_listener_factory( transport: SlackTransport, db_path: Path, enqueue_resume: Callable[[Any], None], - new_task_callback: "Callable[[str, str], str] | None" = None, + new_task_callback: "Callable[[str, str, str], str] | None" = None, ) -> Any: """Build the live :class:`SlackListener` from the coordinator's seams (D-1). @@ -170,13 +170,30 @@ def default_slack_listener_factory( ``new_task_callback`` (WS2) is forwarded to the listener so an allowlisted ``/new-task`` command starts a task. Left ``None``, the listener ignores ``/new-task`` (its built-in default) — the ``/new-task`` path still runs - AFTER the AUTHZ-01 owner check regardless. + AFTER the AUTHZ-01 owner check regardless. The listener posts the root + "📥 Task received" ack through the transport's own poster (no extra wiring). + + A best-effort 👍 ``reactor`` is built from ``SLACK_BOT_TOKEN`` so the listener + can acknowledge inbound answers with a reaction (requires the + ``reactions:write`` scope). If the SDK/token is unavailable the reactor is + left ``None`` (no reaction attempted) — it never blocks listener startup. Imported lazily for the same import-hygiene reason as the clarifier / planner factories (the listener pulls the transport + responder leaves). """ from agent_team.transport.slack_listener import SlackListener + reactor: Any = None + try: + from agent_team.transport.slack_live import build_slack_reactor + + reactor = build_slack_reactor() + except Exception: # noqa: BLE001 - no SDK/token -> run without 👍 reactions + _LOG.info( + "Slack 👍 reactor unavailable (no SDK/token); inbound answers will " + "not be reaction-acknowledged" + ) + return SlackListener( transport, db_path, @@ -184,6 +201,7 @@ def default_slack_listener_factory( app_token=os.environ.get("SLACK_APP_TOKEN") or None, bot_token=os.environ.get("SLACK_BOT_TOKEN") or None, new_task_callback=new_task_callback, + reactor=reactor, ) @@ -449,8 +467,8 @@ class Coordinator: deadline_window: timedelta | None = None, 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, + new_task_callback: "Callable[[str, str, str], str] | None" = None, + notify: "Callable[..., None] | None" = None, ) -> None: self._db_path = Path(db_path) self._transport = transport @@ -505,14 +523,17 @@ class Coordinator: # ------------------------------------------------------------------ # def set_new_task_callback( - self, callback: "Callable[[str, str], str] | None" + self, callback: "Callable[[str, str, str], str] | None" ) -> None: """Wire the WS2 ``/new-task`` handler (call before :meth:`serve`). Avoids the constructor chicken-and-egg of referencing the coordinator's own ``start_task`` at build time: the serve path constructs the - coordinator, then sets ``lambda text, source: self.start_task(...)``. Must - be set before :meth:`_maybe_start_slack_listener` runs (i.e. before + coordinator, then sets + ``lambda text, source, root_ts: self.start_task(...)``. The ``root_ts`` + (the listener's "📥 Task received" ack ts) is forwarded into + ``start_task`` so the task threads under it (one-thread-per-task). Must be + set before :meth:`_maybe_start_slack_listener` runs (i.e. before :meth:`serve`); it is read when the listener is built. """ self._new_task_callback = callback @@ -621,7 +642,9 @@ class Coordinator: # Intake. # ------------------------------------------------------------------ # - def start_task(self, *, task_text: str, transport_name: str) -> str: + def start_task( + self, *, task_text: str, transport_name: str, slack_thread_ts: str = "" + ) -> str: """Start one task: run to the first human gate, then notify (§3.3, §3.3.1). Runs :func:`agent_team.graph.start_task` to the first clarifier @@ -639,6 +662,15 @@ class Coordinator: post-suspend ``update_state`` would clear the pending interrupt and break the human gate). ``intake_node`` returns only a partial state (status/phase), so the seeded ``task`` channel persists into CLARIFY. + + **One-thread-per-task (``slack_thread_ts``).** When the task originates + from a ``/new-task`` slash command, the listener has already posted a + root "📥 Task received" message and passes its ``ts`` here. It is seeded + into the graph state (so every later turn can recover it) AND forwarded + to :func:`notify_question` as ``thread_ts``, so the FIRST clarifier + question posts as a threaded reply under that root and its ledger + ``channel_ref`` becomes the root ``ts``. Empty (the default) for a task + with no root post — the question posts top-level exactly as before. """ if self._graph is None: raise RuntimeError("Coordinator.start_task called before setup()") @@ -652,7 +684,10 @@ class Coordinator: task_text[:200].replace("\n", "\\n").replace("\r", "\\r"), ) thread_id, _state = graph_mod.start_task( - self._graph, transport=transport_name, task=task_text + self._graph, + transport=transport_name, + task=task_text, + slack_thread_ts=slack_thread_ts, ) question = graph_mod.pending_question(self._graph, thread_id=thread_id) @@ -673,6 +708,7 @@ class Coordinator: self._transport, question_set, deadline=deadline, + thread_ts=slack_thread_ts or None, ) finally: conn.close() @@ -747,15 +783,32 @@ class Coordinator: # Maintenance tick (deadline policy + drain). # ------------------------------------------------------------------ # - def _emit(self, message: str) -> None: + def _emit(self, message: str, *, thread_ts: str | None = None) -> 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. + + ``thread_ts`` (one-thread-per-task) — when set, the milestone is posted + as a threaded reply under that task's root "📥 Task received" message so + every notification for a task lands in its one Slack thread. The notify + sink accepts an optional keyword ``thread_ts``; a sink that does not + (an older/simpler sink) is called positionally so the threading hint is + a no-op rather than a crash — the call is tried with ``thread_ts`` first + and falls back to the bare message on a ``TypeError``. """ if self._notify is None: return try: + if thread_ts: + try: + self._notify(message, thread_ts=thread_ts) + return + except TypeError: + # The sink does not accept thread_ts; degrade to a top-level + # post rather than dropping the milestone entirely. + self._notify(message) + return 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) @@ -792,6 +845,11 @@ class Coordinator: desc = desc[:90] + "…" label = f'"{desc}" (`{short}`)' + # One-thread-per-task: the root "📥 Task received" message ts. When + # set, every follow-up question AND every lifecycle milestone for this + # task threads under it. Empty for a non-/new-task origin (top-level). + root_ts = str(values.get("slack_thread_ts") or "") or None + try: question = graph_mod.pending_question(self._graph, thread_id=thread_id) except Exception: # noqa: BLE001 - never let a status check break the loop @@ -800,7 +858,8 @@ class Coordinator: 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. + # 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: @@ -810,6 +869,7 @@ class Coordinator: question["question_set"], deadline=question.get("deadline") or self._default_deadline(), + thread_ts=root_ts, ) finally: conn.close() @@ -819,7 +879,8 @@ class Coordinator: ) self._emit( f"❓ {label} — needs more input. A new clarifying question was " - "posted above; reply in its thread." + "posted above; reply in its thread.", + thread_ts=root_ts, ) continue @@ -844,10 +905,14 @@ class Coordinator: f"• Reached phase: {phase}\n" f"• What's blocking it: {blocker}\n" "• Needs your review — the plan could not be auto-approved. " - "Re-assign with more detail, or adjust the requirement to unblock." + "Re-assign with more detail, or adjust the requirement to unblock.", + thread_ts=root_ts, ) else: - self._emit(f"✅ {label} — plan ready for review (phase: {phase}).") + self._emit( + f"✅ {label} — plan ready for review (phase: {phase}).", + thread_ts=root_ts, + ) @staticmethod def _summarize_blocker(values: "dict[str, Any]") -> str: @@ -995,11 +1060,20 @@ class Coordinator: or question.get("deadline") or self._default_deadline() ) + # One-thread-per-task: re-thread the re-posted question under the + # task's root message so its channel_ref is the root ts again and + # the inbound reply still maps back. Read the root ts from the live + # state (robust across the stub + live clarifier node), falling + # back to the interrupt payload. + root_ts = self._task_root_ts(thread_id) or question.get( + "slack_thread_ts" + ) responder_mod.notify_question( conn, self._transport, question_set, deadline=deadline, + thread_ts=str(root_ts) if root_ts else None, ) finally: conn.close() @@ -1321,6 +1395,22 @@ class Coordinator: # Internals. # ------------------------------------------------------------------ # + def _task_root_ts(self, thread_id: str) -> str | None: + """Return the task's Slack root-message ts (one-thread-per-task), or None. + + Reads ``slack_thread_ts`` off the live checkpointed state so milestone / + re-delivery posts can thread under the task's root "📥 Task received" + message. Best-effort: any read failure (or an empty value) yields + ``None``, so the caller falls back to a top-level post and a status read + never breaks the loop. + """ + try: + snap = self._graph.get_state(graph_mod.thread_config(thread_id)) + values = getattr(snap, "values", {}) or {} + except Exception: # noqa: BLE001 - a state read must not break delivery + return None + return str(values.get("slack_thread_ts") or "") or None + def _default_deadline(self) -> str: """Compute a fresh ISO deadline from the configured window (§3.3.1).""" from datetime import datetime, timezone diff --git a/agent-team/agent_team/graph.py b/agent-team/agent_team/graph.py index 540d547..a724ec5 100644 --- a/agent-team/agent_team/graph.py +++ b/agent-team/agent_team/graph.py @@ -224,6 +224,10 @@ def clarify_node(state: PipelineState) -> PipelineState: turn = len(history) thread_id = state.get("thread_id", "") transport = state.get("transport", "") + # The Slack root-message ts (one-thread-per-task). Empty for non-/new-task + # origins; surfaced in the interrupt payload so the responder can thread the + # question post under the root message (and use it as the channel_ref). + slack_thread_ts = state.get("slack_thread_ts", "") # Stable across the resume replay of this node (see _question_id_for): the # id delivered at suspend == the ledger key == the qa_history entry. @@ -248,6 +252,7 @@ def clarify_node(state: PipelineState) -> PipelineState: "question_set": question_set, "transport": transport, "deadline": deadline, + "slack_thread_ts": slack_thread_ts, } ) @@ -592,6 +597,7 @@ def start_task( thread_id: str | None = None, transport: str = "", task: str = "", + slack_thread_ts: str = "", ) -> tuple[str, PipelineState]: """Start a new pipeline task and run it up to the first human gate (§3.3). @@ -608,6 +614,13 @@ def start_task( empty ``task`` (the default) seeds no description — the clarifier then asks for one. + ``slack_thread_ts`` is the Slack root-message ``ts`` for a ``/new-task`` task + (the "📥 Task received" ack post) — when set, every clarifier question and + lifecycle notification for this task threads under it (one-thread-per-task). + Like ``task`` it is seeded into the initial invoke and persists through + INTAKE into CLARIFY (``intake_node`` returns only a partial state). Empty + (the default) for a task with no root post — posts are top-level as before. + The graph MUST be compiled with a checkpointer for the suspend to persist; an uncheckpointed graph would run straight through without honouring the interrupt. @@ -619,6 +632,7 @@ def start_task( status=TaskStatus.ACTIVE.value, current_phase=_phase_value(Phase.INTAKE), task=task, + slack_thread_ts=slack_thread_ts, qa_history=[], transport=transport, created_at=now, diff --git a/agent-team/agent_team/responder.py b/agent-team/agent_team/responder.py index 96333ac..997004c 100644 --- a/agent-team/agent_team/responder.py +++ b/agent-team/agent_team/responder.py @@ -156,6 +156,7 @@ def notify_question( *, deadline: str, posted_at: str | None = None, + thread_ts: str | None = None, ) -> str | None: """Deliver a question-set: write the ledger row ``open`` first, then post. @@ -177,6 +178,16 @@ def notify_question( The row is inserted with the question's identity (``thread_id``, ``turn``, ``transport`` name) so the turn guard and reconcile can act on it. + + ``thread_ts`` (one-thread-per-task, Slack) — when set, the question is posted + as a THREADED REPLY under that root message ``ts`` (the task's + "📥 Task received" ack post) AND the row's durable ``channel_ref`` is set to + that SAME root ``ts`` (NOT the posted reply's own ``ts``). This is what makes + the answer-mapping unchanged: a human reply in the root thread carries + ``thread_ts == root_ts``, and ``find_open_question_by_channel_ref(thread_ts)`` + resolves it to this task's currently-open question. When ``None`` the + question is posted top-level and the ``channel_ref`` is the posted message's + own ``ts`` exactly as before. """ stamp = posted_at or _utc_now_iso() transport_name = type(transport).__name__ @@ -197,19 +208,29 @@ def notify_question( ) # 2. Side-effecting post. A failure here is recoverable (row stays open, - # no ref) — do NOT let it bubble up and lose the durable row. + # no ref) — do NOT let it bubble up and lose the durable row. ``thread_ts`` + # is only forwarded when set, so non-threading transports keep their + # existing call shape. try: - channel_ref = transport.post_question( + post_kwargs: dict[str, Any] = {} + if thread_ts: + post_kwargs["thread_ts"] = thread_ts + posted_ref = transport.post_question( thread_id=question_set.thread_id, question_id=question_set.question_id, turn=question_set.turn, question_set=question_set, deadline=deadline, + **post_kwargs, ) except Exception: return None - # 3. Persist the ref so reconcile/recovery can act on the post. + # 3. Persist the ref so reconcile/recovery can act on the post. When the post + # threaded under a root message, the durable channel_ref is the ROOT ts + # (so an inbound reply's thread_ts maps back to this question via + # find_open_question_by_channel_ref), NOT the posted reply's own ts. + channel_ref = thread_ts if thread_ts else posted_ref conn.execute( "UPDATE pending_questions SET channel_ref=? WHERE question_id=?", (channel_ref, question_set.question_id), diff --git a/agent-team/agent_team/task_model.py b/agent-team/agent_team/task_model.py index 9408ed7..7a26a03 100644 --- a/agent-team/agent_team/task_model.py +++ b/agent-team/agent_team/task_model.py @@ -89,6 +89,9 @@ class TaskRecord: current_phase: Phase # Intake task description (mirrors PipelineState.task). task: str = "" + # Slack root-message ts for one-thread-per-task (mirrors + # PipelineState.slack_thread_ts). Empty for non-/new-task origins. + slack_thread_ts: str = "" qa_history: list[Any] = field(default_factory=list) plan: dict[str, Any] | None = None review_verdicts: list[Any] = field(default_factory=list) @@ -114,6 +117,12 @@ class PipelineState(TypedDict, total=False): # Seeded by graph.start_task and read by the clarifier/planner; a first-class # channel so the seeded value persists across node transitions. task: str + # The Slack root-message ``ts`` for a /new-task task (the "📥 Task received" + # ack post). All of the task's clarifier questions and lifecycle milestone + # notifications thread under this ``ts`` so one task maps to one Slack thread. + # Empty/absent for a task that did not originate from /new-task (no root post), + # in which case posts are top-level exactly as before. + slack_thread_ts: str qa_history: list[Any] plan: dict[str, Any] | None review_verdicts: list[Any] @@ -140,6 +149,7 @@ def task_from_dict(data: dict[str, Any]) -> TaskRecord: status=TaskStatus(data["status"]), current_phase=Phase(data["current_phase"]), task=data.get("task", ""), + slack_thread_ts=data.get("slack_thread_ts", ""), qa_history=list(data.get("qa_history", [])), plan=data.get("plan"), review_verdicts=list(data.get("review_verdicts", [])), diff --git a/agent-team/agent_team/transport/base.py b/agent-team/agent_team/transport/base.py index 015cfda..2779816 100644 --- a/agent-team/agent_team/transport/base.py +++ b/agent-team/agent_team/transport/base.py @@ -84,6 +84,7 @@ class Transport(ABC): turn: int, question_set: QuestionSet, deadline: str, + thread_ts: str | None = None, ) -> str: """Deliver ``question_set`` and return its ``channel_ref``. @@ -93,6 +94,12 @@ class Transport(ABC): ``channel_ref`` is the transport's locator for the post (Slack message ``ts`` / issue-comment id / Claude session id) and is stored on the ledger row so reconcile/recovery can act on it (§3.3.1). + + ``thread_ts`` is an OPTIONAL transport-specific threading hint (Slack's + one-thread-per-task: post the message as a reply under that root ``ts``). + Transports without native threading may ignore it. The responder only + forwards it when set, so a transport that does not accept it is never + called with it. """ raise NotImplementedError diff --git a/agent-team/agent_team/transport/slack_adapter.py b/agent-team/agent_team/transport/slack_adapter.py index 5c0f7c3..417c33b 100644 --- a/agent-team/agent_team/transport/slack_adapter.py +++ b/agent-team/agent_team/transport/slack_adapter.py @@ -217,6 +217,19 @@ class SlackTransport(Transport): self.channel = channel self._poster: SlackPoster = poster if poster is not None else _default_poster + @property + def poster(self) -> SlackPoster: + """The injected network seam (the ``chat.postMessage`` callable). + + Exposed read-only so collaborators sharing this transport (e.g. the + inbound :class:`~agent_team.transport.slack_listener.SlackListener`) can + post auxiliary messages — the root "📥 Task received" ack and the 👍 + reaction-bearing posts — through the SAME poster the question delivery + uses, instead of constructing a second client. The foundation default + still refuses the network (no poster configured). + """ + return self._poster + def post_question( self, *, @@ -225,12 +238,23 @@ class SlackTransport(Transport): turn: int, question_set: QuestionSet, deadline: str, + thread_ts: str | None = None, ) -> str: """Render + post the question-set; return the Slack ``ts`` channel_ref. Embeds ``question_id`` in the message ``callback_id`` so an inbound answer maps back (§3.3.1). On any poster failure raises :class:`SlackPostError` so the ledger row stays ``open`` for reconcile. + + ``thread_ts`` (one-thread-per-task) — when set, the message is posted as + a THREADED REPLY under that root ``ts`` (the task's "📥 Task received" + ack post), so every clarifier question for a task lands in one Slack + thread. When ``None`` (the default, e.g. a task that did not originate + from ``/new-task``) the message is posted top-level exactly as before. + The returned value is still the POSTED message's own ``ts``; the caller + (:func:`agent_team.responder.notify_question`) is what records the + durable ``channel_ref`` (it uses the root ``thread_ts`` when threading so + an inbound reply's ``thread_ts`` maps back to this question). """ blocks = build_question_blocks(question_set, deadline) message: dict[str, Any] = { @@ -250,6 +274,11 @@ class SlackTransport(Transport): }, }, } + # Thread under the task's root message when one exists (one thread per + # task). Only set the key when non-empty so the top-level-post behavior + # is byte-identical for non-/new-task origins. + if thread_ts: + message["thread_ts"] = thread_ts try: response = self._poster(message) diff --git a/agent-team/agent_team/transport/slack_listener.py b/agent-team/agent_team/transport/slack_listener.py index 7df310d..685b86d 100644 --- a/agent-team/agent_team/transport/slack_listener.py +++ b/agent-team/agent_team/transport/slack_listener.py @@ -82,17 +82,29 @@ from collections.abc import Callable from agent_team.db.schema import connect, find_open_question_by_channel_ref from agent_team.responder import AnswerOutcome, EnqueueResume, submit_answer -from agent_team.transport.slack_adapter import SlackTransport +from agent_team.transport.slack_adapter import SlackPoster, SlackTransport __all__ = [ "NewTaskCallback", + "Reactor", "SlackListener", ] -# Injectable callback for /new-task slash commands: receives (task_text, transport) -# and returns the minted thread_id. Injected at coordinator startup so the -# listener is testable with no coordinator and no graph. -NewTaskCallback = Callable[[str, str], str] +# Injectable callback for /new-task slash commands: receives +# (task_text, transport, slack_thread_ts) and returns the minted thread_id. +# ``slack_thread_ts`` is the root "📥 Task received" message ts the listener +# posted before starting the task (one-thread-per-task); the callback seeds it +# into the graph state so every later question/notification threads under it. +# Injected at coordinator startup so the listener is testable with no +# coordinator and no graph. +NewTaskCallback = Callable[[str, str, str], str] + +# A best-effort 👍-reaction adder: given (channel, message_ts), add a reaction so +# the human sees the machine received the inbound message. Returns nothing; any +# failure (missing scope, deleted message, transport error) must be tolerated by +# the caller. Injected so the listener is testable with a fake recorder and no +# live WebClient. +Reactor = Callable[[str, str], None] # The Slack slash command that starts a new pipeline task. _NEW_TASK_COMMAND = "/new-task" @@ -162,10 +174,20 @@ class SlackListener: answer; :meth:`serve` sources it from ``AGENT_TEAM_SLACK_OWNER_IDS`` when not injected. * ``new_task_callback`` — optional :data:`NewTaskCallback`; when set, the - listener handles ``/new-task `` slash commands by calling it - with ``(task_text, "slack")`` and returning ``None`` (the task is started; - the owner will receive clarifying questions via the transport). When - ``None`` (the default), ``/new-task`` commands are ignored. + listener handles ``/new-task `` slash commands by posting a + root "📥 Task received" ack message (one-thread-per-task) and calling it + with ``(task_text, "slack", root_ts)``, returning ``None`` (the task is + started; the owner will receive clarifying questions threaded under the + root). When ``None`` (the default), ``/new-task`` commands are ignored. + * ``poster`` — optional :data:`~agent_team.transport.slack_adapter.SlackPoster` + used ONLY to post the root "📥 Task received" ack for ``/new-task``. + Defaults to the injected ``transport``'s own poster so the ack goes through + the same client as the questions. Never used for answers. + * ``reactor`` — optional :data:`Reactor`; when set, the listener adds a 👍 + reaction to an inbound message it acted on (a thread-reply answer) AFTER + the AUTHZ-01 owner check passes, best-effort. ``None`` (the default) means + no reaction is attempted. Requires the ``reactions:write`` bot scope; until + that is granted the reactor silently no-ops, which the listener tolerates. The listener never resumes the graph; it only normalizes, submits, and enqueues. See the module SECURITY note for the trust boundary. @@ -181,6 +203,8 @@ class SlackListener: bot_token: str | None = None, owner_ids: set[str] | None = None, new_task_callback: NewTaskCallback | None = None, + poster: SlackPoster | None = None, + reactor: Reactor | None = None, ) -> None: self._transport = transport self._db_path = Path(db_path) @@ -192,6 +216,11 @@ class SlackListener: self._owner_ids: set[str] = set(owner_ids) if owner_ids else set() # Injected /new-task callback (opt-in). None = ignore new-task commands. self._new_task_callback = new_task_callback + # Poster for the root "📥 Task received" ack. Defaults to the transport's + # own poster so the ack uses the same client as the question posts. + self._poster: SlackPoster = poster if poster is not None else transport.poster + # 👍-reaction adder (opt-in). None = no reaction attempted. + self._reactor = reactor # The live Socket Mode handler, retained by :meth:`serve` so :meth:`close` # can stop it cleanly on daemon shutdown. ``None`` until ``serve`` opens # the socket. @@ -249,9 +278,17 @@ class SlackListener: # slash command is never misrouted as an answer-to-a-question. Auth has # already cleared above (AUTHZ-01), so only allowlisted owners can start # tasks. The callback is opt-in; if not injected, /new-task is ignored. + # (No 👍 reaction here: a slash command has no reactable message; its ack + # is the "📥 Task received" root post instead.) if _is_new_task_command(raw_payload): return self._handle_new_task_command(raw_payload) + # 👍-acknowledge the inbound message the machine is acting on (an answer + # in a task thread). Runs AFTER AUTHZ-01 (a non-owner message above + # already returned None, so this never reacts to an unauthorized sender) + # and is best-effort: a reaction failure must never break handle_event. + self._maybe_react(raw_payload) + # ``submit_answer`` calls ``transport.parse_answer`` internally, which # raises ValueError when no question_id is recoverable. A real free-text # thread reply carries no callback_id / question_id / metadata, so its @@ -315,12 +352,21 @@ class SlackListener: Extracts the task description from ``text`` (the words after the command name). If no ``new_task_callback`` is configured, logs and returns - ``None`` (ignore). Otherwise calls ``new_task_callback(task_text, "slack")`` - and returns ``None`` (the task is started; the owner receives clarifying - questions via the transport; there is no ``AnswerOutcome`` to return here). + ``None`` (ignore). Otherwise: - Any exception raised by the callback is caught and logged; the listen - loop stays alive. + 1. Posts an immediate ROOT "📥 Task received" ack message to the channel + (one-thread-per-task) and captures its ``ts`` (``root_ts``). This is + the instant acknowledgement (a slash command has no reactable message, + so this post IS its 👍). A post failure degrades to ``root_ts=""`` so + the task still starts (its questions then post top-level). + 2. Calls ``new_task_callback(task_text, "slack", root_ts)`` so the task is + seeded with the root ts and every clarifier question + lifecycle + notification threads under it. + + Returns ``None`` (the task is started; the owner receives clarifying + questions via the transport; there is no ``AnswerOutcome`` here). Any + exception raised by the callback is caught and logged; the listen loop + stays alive. """ if self._new_task_callback is None: _LOG.debug( @@ -333,11 +379,16 @@ class SlackListener: _LOG.info("/new-task received with empty description; ignoring") return None + # 1. Instant root ack (one-thread-per-task). Best-effort: a failed post + # yields root_ts="" so the task still starts (top-level questions). + root_ts = self._post_task_received(text) + try: - thread_id = self._new_task_callback(text, "slack") + thread_id = self._new_task_callback(text, "slack", root_ts) _LOG.info( - "new task started via /new-task: thread_id=%s task=%r", + "new task started via /new-task: thread_id=%s root_ts=%s task=%r", thread_id, + root_ts or "(none)", text[:80], ) except Exception: # noqa: BLE001 - keep the listen loop alive @@ -348,6 +399,76 @@ class SlackListener: ) return None + def _post_task_received(self, description: str) -> str: + """Post the root "📥 Task received" ack to the channel; return its ``ts``. + + The instant acknowledgement for a ``/new-task`` command and the anchor for + one-thread-per-task: every clarifier question and lifecycle notification + threads under the returned ``ts``. Posts through the shared poster (the + same client the question delivery uses) to the transport's channel. + + Best-effort: any failure (no poster configured, transport error, missing + ``ts`` in the response) is swallowed and an empty string returned, so a + post problem never blocks task start — the task simply runs with + top-level (un-threaded) questions. + """ + try: + response = self._poster( + { + "channel": self._transport.channel, + "text": ( + f'📥 Task received: "{description}" ' + "— starting (clarifying first)…" + ), + } + ) + except Exception: # noqa: BLE001 - a failed ack must not block task start + _LOG.warning( + "failed to post '📥 Task received' root ack; task starts un-threaded", + exc_info=True, + ) + return "" + if not isinstance(response, Mapping): + return "" + ts = response.get("ts") + if not ts: + message = response.get("message") + if isinstance(message, Mapping): + ts = message.get("ts") + return str(ts) if ts else "" + + def _maybe_react(self, raw_payload: Mapping[str, Any]) -> None: + """Add a 👍 reaction to the inbound message the machine is acting on. + + Best-effort acknowledgement that the inbound answer was received. Only + fires when a ``reactor`` is configured and the payload is an Events API + message/mention carrying a channel + message ``ts`` (the reactable + thread-reply answer). MUST be called only AFTER AUTHZ-01 has passed (the + caller guarantees this), so a non-owner message is never reacted to. + + Wrapped end-to-end: a missing scope (``reactions:write`` not yet granted), + a deleted message, or any transport error is swallowed — a reaction + failure must never break :meth:`handle_event` or the listen loop. + """ + if self._reactor is None: + return + event = _inner_event(raw_payload) + # Only react to a real inbound message/mention (the thread-reply answer + # shape). Interactive/slash payloads have no reactable message ts here. + channel = event.get("channel") + ts = event.get("ts") + if not channel or not ts: + return + try: + self._reactor(str(channel), str(ts)) + except Exception: # noqa: BLE001 - a reaction failure must never break handling + _LOG.debug( + "👍 reaction add failed (channel=%s ts=%s); ignoring " + "(reactions:write may not be granted yet)", + channel, + ts, + ) + def _is_authorized(self, raw_payload: Mapping[str, Any]) -> bool: """Return ``True`` iff the payload's sender is an allowlisted owner. diff --git a/agent-team/agent_team/transport/slack_live.py b/agent-team/agent_team/transport/slack_live.py index 9f9868e..1a0b64d 100644 --- a/agent-team/agent_team/transport/slack_live.py +++ b/agent-team/agent_team/transport/slack_live.py @@ -38,7 +38,7 @@ do the question-id round-trip the adapter relies on. from __future__ import annotations import os -from collections.abc import Mapping +from collections.abc import Callable, Mapping from typing import Any from agent_team.transport.slack_adapter import SlackPoster, SlackTransport @@ -46,12 +46,15 @@ from agent_team.transport.slack_adapter import SlackPoster, SlackTransport __all__ = [ "build_live_slack_transport", "build_slack_poster", + "build_slack_reactor", ] # Top-level ``chat.postMessage`` keyword arguments the live poster forwards. # ``callback_id`` is deliberately excluded: it is not a postMessage parameter, # and the durable inbound key lives in ``metadata.event_payload`` instead. -_POST_MESSAGE_KEYS = ("channel", "text", "blocks", "metadata") +# ``thread_ts`` IS a postMessage parameter (one-thread-per-task threading) and is +# forwarded when present so a question/notification posts as a threaded reply. +_POST_MESSAGE_KEYS = ("channel", "text", "blocks", "metadata", "thread_ts") def build_slack_poster(token: str | None = None, *, client: Any = None) -> SlackPoster: @@ -86,6 +89,39 @@ def build_slack_poster(token: str | None = None, *, client: Any = None) -> Slack return _poster +def build_slack_reactor( + token: str | None = None, *, client: Any = None +) -> "Callable[[str, str], None]": + """Build a live ``slack_sdk``-backed 👍-reaction adder (one-thread-per-task UX). + + The returned ``(channel, ts) -> None`` callable performs a Slack + ``reactions.add`` (emoji ``thumbsup``) on the message at ``(channel, ts)`` so + a human sees the machine received their inbound answer. It is wired into the + inbound :class:`~agent_team.transport.slack_listener.SlackListener` (which + only calls it AFTER the AUTHZ-01 owner check passes) and is invoked + best-effort — the listener swallows any failure. + + Requires the ``reactions:write`` bot scope. Until that scope is granted (the + manifest re-applied + the app reinstalled) ``reactions.add`` fails with a + ``missing_scope`` error; this reactor lets that propagate to the listener, + which swallows it, so the reaction silently no-ops rather than breaking + answer handling. + + ``client`` (optional) injects a pre-built client for testability; any object + exposing ``reactions_add(**kwargs)`` works. When omitted, a + ``slack_sdk.WebClient`` is constructed lazily from ``token`` (falling back to + ``SLACK_BOT_TOKEN``); the deferred-import / missing-token semantics match + :func:`build_slack_poster`. + """ + if client is None: + client = _build_web_client(token) + + def _reactor(channel: str, ts: str) -> None: + client.reactions_add(channel=channel, timestamp=ts, name="thumbsup") + + return _reactor + + def build_live_slack_transport( channel: str, token: str | None = None, *, client: Any = None ) -> SlackTransport: diff --git a/agent-team/run-team.py b/agent-team/run-team.py index 57f9b09..d699217 100644 --- a/agent-team/run-team.py +++ b/agent-team/run-team.py @@ -616,8 +616,10 @@ def _build_coordinator(args: argparse.Namespace) -> Any: # before serve() builds the listener. AUTHZ-01 (owner allowlist) gates this # upstream in the listener; the source label is always "slack". coordinator.set_new_task_callback( - lambda task_text, _source: coordinator.start_task( - task_text=task_text, transport_name="slack" + lambda task_text, _source, slack_thread_ts: coordinator.start_task( + task_text=task_text, + transport_name="slack", + slack_thread_ts=slack_thread_ts, ) ) return coordinator diff --git a/agent-team/tests/test_coordinator.py b/agent-team/tests/test_coordinator.py index 66a6d2b..6de75fc 100644 --- a/agent-team/tests/test_coordinator.py +++ b/agent-team/tests/test_coordinator.py @@ -53,6 +53,9 @@ class FakeTransport(Transport): def __init__(self, *, fail_post: bool = False) -> None: self.posted: list[QuestionSet] = [] + # Records the thread_ts each post was threaded under (None = top-level), + # so one-thread-per-task wiring can be asserted. + self.thread_tss: list[str | None] = [] self.fail_post = fail_post def post_question( @@ -63,10 +66,12 @@ class FakeTransport(Transport): turn: int, question_set: QuestionSet, deadline: str, + thread_ts: str | None = None, ) -> str: if self.fail_post: raise RuntimeError("simulated transport post failure") self.posted.append(question_set) + self.thread_tss.append(thread_ts) return f"fake:{question_id}" def parse_answer(self, raw: Any) -> tuple[str, Any, str]: @@ -201,6 +206,53 @@ def test_start_task_suspends_and_writes_open_ledger_row(db_path: Path) -> None: assert len(transport.posted) == 1 +def test_start_task_threads_first_question_under_root(db_path: Path) -> None: + """One-thread-per-task: a /new-task root ts threads the first question + ref. + + When start_task is given ``slack_thread_ts`` (the listener's "📥 Task + received" ack ts), the first clarifier question posts threaded under it AND + the ledger ``channel_ref`` is that ROOT ts — so the human's reply in the + thread (thread_ts == root) maps back to this open question with NO change to + the answer-mapping logic. + """ + transport = FakeTransport() + coord = _make_coordinator(db_path, transport=transport) + coord.setup() + + root_ts = "1700000000.ROOT" + thread_id = coord.start_task( + task_text="build a thing", + transport_name="slack", + slack_thread_ts=root_ts, + ) + assert thread_id + + # The question post threaded under the root. + assert transport.thread_tss == [root_ts] + # The ledger channel_ref is the ROOT ts (so an inbound reply's thread_ts maps + # back via find_open_question_by_channel_ref), not the posted reply's ref. + row = _only_open_row(db_path) + assert row["channel_ref"] == root_ts + # And the root ts is seeded on the graph state for later turns/notifications. + state = graph_mod.get_pipeline_state(coord.graph, thread_id=thread_id) + assert state.get("slack_thread_ts") == root_ts + + +def test_start_task_without_root_posts_top_level(db_path: Path) -> None: + """No slack_thread_ts (e.g. a non-/new-task origin): top-level, ref = post ref.""" + transport = FakeTransport() + coord = _make_coordinator(db_path, transport=transport) + coord.setup() + + thread_id = coord.start_task(task_text="x", transport_name="github") + + assert transport.thread_tss == [None] + row = _only_open_row(db_path) + assert row["channel_ref"] == f"fake:{row['question_id']}" + state = graph_mod.get_pipeline_state(coord.graph, thread_id=thread_id) + assert state.get("slack_thread_ts") == "" + + def test_start_task_lost_post_leaves_open_row_without_ref(db_path: Path) -> None: # A failed transport post is recoverable: the row stays open with no ref. transport = FakeTransport(fail_post=True) @@ -953,7 +1005,9 @@ def test_followups_posts_new_question_and_emits_needs_input( monkeypatch.setattr( coord_mod.responder_mod, "notify_question", - lambda conn, transport, qset, *, deadline: posted.append((qset, deadline)), + lambda conn, transport, qset, *, deadline, thread_ts=None: posted.append( + (qset, deadline, thread_ts) + ), ) coord._post_resume_followups([_resume_result("abc12345deadbeef")]) @@ -962,6 +1016,51 @@ def test_followups_posts_new_question_and_emits_needs_input( assert any("abc12345" in m for m in msgs) +def test_followups_thread_under_task_root_ts(db_path: Path, monkeypatch: Any) -> None: + """A multi-turn follow-up question + milestone thread under the task root ts. + + The task's ``slack_thread_ts`` (seeded by start_task) is read off the live + state and forwarded as ``thread_ts`` to both the follow-up notify_question + and the milestone _emit, so one-thread-per-task holds across turns. + """ + from agent_team import coordinator as coord_mod + + posted: list[Any] = [] + emitted: list[tuple[str, Any]] = [] + coord = _make_coordinator(db_path) + # A notify sink that accepts the optional thread_ts kwarg. + coord._notify = lambda message, *, thread_ts=None: emitted.append( + (message, thread_ts) + ) + coord.setup() + + # A real task carrying a root ts on its state. + root_ts = "1700000000.ROOT" + thread_id = coord.start_task( + task_text="ship it", transport_name="slack", slack_thread_ts=root_ts + ) + + 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, thread_ts=None: posted.append( + thread_ts + ), + ) + + coord._post_resume_followups([_resume_result(thread_id)]) + + # The follow-up question threaded under the root ts. + assert posted == [root_ts] + # The "needs more input" milestone also threaded under the root ts. + assert any(ts == root_ts for _msg, ts in emitted) + + 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 diff --git a/agent-team/tests/test_graph.py b/agent-team/tests/test_graph.py index 48d0b97..b4b58be 100644 --- a/agent-team/tests/test_graph.py +++ b/agent-team/tests/test_graph.py @@ -147,6 +147,28 @@ def test_start_task_seeds_task_description_into_state(compiled) -> None: assert state.get("task") == "build a login form" +def test_start_task_seeds_slack_thread_ts_into_state(compiled) -> None: + # One-thread-per-task: the root "📥 Task received" message ts must reach the + # graph state so every later question/notification threads under it. Like + # ``task`` it is seeded into the initial invoke and persists through INTAKE + # into the suspended CLARIFY snapshot. + _thread_id, state = start_task( + compiled, transport="slack", slack_thread_ts="1700000000.ROOT" + ) + assert state.get("slack_thread_ts") == "1700000000.ROOT" + # It is also surfaced on the pending interrupt payload so the responder can + # thread the question post. + payload = pending_question(compiled, thread_id=_thread_id) + assert payload["slack_thread_ts"] == "1700000000.ROOT" + + +def test_start_task_default_slack_thread_ts_is_empty(compiled) -> None: + # A task with no root post (e.g. a GitHub-issue origin): slack_thread_ts is + # empty so questions post top-level exactly as before. + _thread_id, state = start_task(compiled, transport="slack") + assert state.get("slack_thread_ts") == "" + + def test_pending_question_carries_foundation_questionset(compiled) -> None: thread_id, _ = start_task(compiled, transport="slack") payload = pending_question(compiled, thread_id=thread_id) diff --git a/agent-team/tests/test_responder.py b/agent-team/tests/test_responder.py index 914b6ae..b541d29 100644 --- a/agent-team/tests/test_responder.py +++ b/agent-team/tests/test_responder.py @@ -71,7 +71,7 @@ class FakeTransport(Transport): self.post_fails = post_fails def post_question( - self, *, thread_id, question_id, turn, question_set, deadline + self, *, thread_id, question_id, turn, question_set, deadline, thread_ts=None ) -> str: if self.post_fails: raise RuntimeError("transport unreachable") @@ -82,6 +82,7 @@ class FakeTransport(Transport): "question_id": question_id, "turn": turn, "deadline": deadline, + "thread_ts": thread_ts, "ref": ref, } ) @@ -197,6 +198,56 @@ def test_notify_lost_post_leaves_open_row_without_ref( assert row["channel_ref"] is None +def test_notify_threads_under_root_and_sets_channel_ref_to_root( + conn: sqlite3.Connection, +) -> None: + """One-thread-per-task: thread_ts is forwarded to post AND becomes channel_ref. + + When ``thread_ts`` (the task's root "📥 Task received" ts) is given, the + question posts as a threaded reply under it, and the durable ``channel_ref`` + is set to that ROOT ts (NOT the posted reply's own ts) — so an inbound reply + whose ``thread_ts == root_ts`` maps back via + ``find_open_question_by_channel_ref``. + """ + transport = FakeTransport() + qs = _question_set() + root_ts = "1700000000.ROOT" + + ref = notify_question( + conn, + transport, + qs, + deadline="2026-06-18T00:00:00+00:00", + thread_ts=root_ts, + ) + + # The post threaded under the root. + assert transport.posts[0]["thread_ts"] == root_ts + # channel_ref is the ROOT ts, not the posted reply ref ("slack-ts-q1"). + assert ref == root_ts + row = _row(conn, "q1") + assert row["channel_ref"] == root_ts + + +def test_notify_without_thread_ts_uses_posted_ref_as_channel_ref( + conn: sqlite3.Connection, +) -> None: + """No thread_ts (the default) preserves the prior behavior exactly. + + The post is top-level (thread_ts None) and the channel_ref is the posted + message's own ts (the FakeTransport ref). + """ + transport = FakeTransport() + + ref = notify_question( + conn, transport, _question_set(), deadline="2026-06-18T00:00:00+00:00" + ) + + assert transport.posts[0]["thread_ts"] is None + assert ref == "slack-ts-q1" + assert _row(conn, "q1")["channel_ref"] == "slack-ts-q1" + + # --------------------------------------------------------------------------- # submit_answer — first-answer-wins (§3.3.1). # --------------------------------------------------------------------------- diff --git a/agent-team/tests/test_slack_adapter.py b/agent-team/tests/test_slack_adapter.py index dc9e892..25376e7 100644 --- a/agent-team/tests/test_slack_adapter.py +++ b/agent-team/tests/test_slack_adapter.py @@ -169,6 +169,46 @@ def test_post_question_embeds_question_id_in_callback_id() -> None: assert sent["metadata"]["event_payload"]["thread_id"] == "thread-1" +def test_post_question_threads_under_thread_ts_when_given() -> None: + """One-thread-per-task: a non-empty thread_ts is forwarded as message['thread_ts']. + + The poster receives ``thread_ts`` so Slack posts the question as a threaded + reply under the task's root "📥 Task received" message. + """ + poster = _RecordingPoster() + t = SlackTransport(channel="C999", poster=poster) + t.post_question( + thread_id="thread-1", + question_id="q-abc", + turn=0, + question_set=_question_set(), + deadline="2026-06-18T00:00:00Z", + thread_ts="1700000000.ROOT", + ) + assert poster.calls[0]["thread_ts"] == "1700000000.ROOT" + + +def test_post_question_omits_thread_ts_by_default() -> None: + """No thread_ts (the default) => top-level post (no 'thread_ts' key).""" + poster = _RecordingPoster() + t = SlackTransport(channel="C999", poster=poster) + t.post_question( + thread_id="thread-1", + question_id="q-abc", + turn=0, + question_set=_question_set(), + deadline="2026-06-18T00:00:00Z", + ) + assert "thread_ts" not in poster.calls[0] + + +def test_poster_property_exposes_injected_poster() -> None: + """The transport exposes its poster so the listener can post the root ack.""" + poster = _RecordingPoster() + t = SlackTransport(channel="C1", poster=poster) + assert t.poster is poster + + def test_post_question_accepts_nested_message_ts() -> None: poster = _RecordingPoster({"ok": True, "message": {"ts": "1700000000.000300"}}) t = SlackTransport(channel="C1", poster=poster) diff --git a/agent-team/tests/test_ws0_ws2_ws4_plugin_slack_hook.py b/agent-team/tests/test_ws0_ws2_ws4_plugin_slack_hook.py index 4f198be..d14cb40 100644 --- a/agent-team/tests/test_ws0_ws2_ws4_plugin_slack_hook.py +++ b/agent-team/tests/test_ws0_ws2_ws4_plugin_slack_hook.py @@ -157,23 +157,25 @@ def test_is_new_task_command_ignores_non_slash() -> None: def test_new_task_calls_callback_with_text_and_transport() -> None: - calls: list[tuple[str, str]] = [] + calls: list[tuple[str, str, str]] = [] - def cb(task: str, transport: str) -> str: - calls.append((task, transport)) + def cb(task: str, transport: str, root_ts: str) -> str: + calls.append((task, transport, root_ts)) return "thread-abc" listener = _make_listener(new_task_callback=cb) result = listener.handle_event(_new_task_payload("Add OAuth to admin portal")) assert result is None # no AnswerOutcome for new-task - assert calls == [("Add OAuth to admin portal", "slack")] + # The default (non-posting) transport poster raises, so root_ts degrades to + # "" — the task still starts, un-threaded. The text + via are forwarded. + assert calls == [("Add OAuth to admin portal", "slack", "")] def test_new_task_unauthorized_sender_rejected() -> None: calls: list[Any] = [] - def cb(task: str, transport: str) -> str: + def cb(task: str, transport: str, root_ts: str) -> str: calls.append(task) return "thread-xyz" @@ -187,7 +189,7 @@ def test_new_task_unauthorized_sender_rejected() -> None: def test_new_task_empty_text_ignored() -> None: calls: list[Any] = [] - def cb(task: str, transport: str) -> str: + def cb(task: str, transport: str, root_ts: str) -> str: calls.append(task) return "thread-123" @@ -205,7 +207,7 @@ def test_new_task_no_callback_configured_is_ignored() -> None: def test_new_task_callback_exception_does_not_crash_listener() -> None: - def boom(task: str, transport: str) -> str: + def boom(task: str, transport: str, root_ts: str) -> str: raise RuntimeError("coordinator exploded") listener = _make_listener(new_task_callback=boom) @@ -222,7 +224,7 @@ def test_other_slash_command_not_intercepted_by_new_task_path( init_db(db_path) calls: list[Any] = [] - def cb(task: str, transport: str) -> str: + def cb(task: str, transport: str, root_ts: str) -> str: calls.append(task) return "thread-xyz" diff --git a/agent-team/tests/test_ws_activation_wiring.py b/agent-team/tests/test_ws_activation_wiring.py index 9988d9c..9e90760 100644 --- a/agent-team/tests/test_ws_activation_wiring.py +++ b/agent-team/tests/test_ws_activation_wiring.py @@ -97,18 +97,27 @@ def test_build_coordinator_wires_new_task_callback_to_start_task( # the real graph. calls: dict[str, Any] = {} - def _fake_start_task(*, task_text: str, transport_name: str) -> str: + def _fake_start_task( + *, task_text: str, transport_name: str, slack_thread_ts: str = "" + ) -> str: calls["task_text"] = task_text calls["transport_name"] = transport_name + calls["slack_thread_ts"] = slack_thread_ts return "thread-xyz" monkeypatch.setattr(coordinator, "start_task", _fake_start_task) cb = coordinator._new_task_callback assert cb is not None - thread_id = cb("fix the flaky test", "slack") + # The 3-arg callback (one-thread-per-task): (task_text, via, root_ts). The + # root_ts is forwarded into start_task so the task threads under the root. + thread_id = cb("fix the flaky test", "slack", "1700000000.000100") assert thread_id == "thread-xyz" - assert calls == {"task_text": "fix the flaky test", "transport_name": "slack"} + assert calls == { + "task_text": "fix the flaky test", + "transport_name": "slack", + "slack_thread_ts": "1700000000.000100", + } def test_set_new_task_callback_overrides() -> None: @@ -124,7 +133,7 @@ def test_set_new_task_callback_overrides() -> None: coord = Coordinator(db_path=":memory:", transport=_T()) assert coord._new_task_callback is None - sentinel = lambda t, s: "tid" # noqa: E731 + sentinel = lambda t, s, r: "tid" # noqa: E731 coord.set_new_task_callback(sentinel) assert coord._new_task_callback is sentinel @@ -150,7 +159,7 @@ def test_slack_listener_factory_forwards_new_task_callback( monkeypatch.setattr(sl_mod, "SlackListener", _FakeListener) - sentinel = lambda t, s: "tid" # noqa: E731 + sentinel = lambda t, s, r: "tid" # noqa: E731 coord_mod.default_slack_listener_factory( transport=MagicMock(), db_path=tmp_path / "x.sqlite", From 9bfb5f1534f68f7b11784e5926ae5ed112926c15 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 15:11:30 -0400 Subject: [PATCH 13/13] =?UTF-8?q?feat(agent-team):=20=F0=9F=91=8D-acknowle?= =?UTF-8?q?dge=20received=20Slack=20answers=20(reactions:write)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit WS Slack-UX Feature 2. When the inbound listener acts on an answer in a task thread, it adds a 👍 reaction to that reply so the human sees the machine received it. - SlackListener gains an optional reactor seam; handle_event reacts to the inbound reply message (channel + event ts) AFTER the AUTHZ-01 owner check passes — a non-owner message is rejected and never reacted to. Best-effort: any reaction failure (notably a missing scope) is swallowed and never breaks handle_event or the listen loop. - /new-task is NOT reacted to (a slash command has no reactable message); its "📥 Task received" root post is the acknowledgement. - build_slack_reactor wraps WebClient.reactions_add(name="thumbsup"); the default listener factory wires it best-effort from SLACK_BOT_TOKEN. - Adds reactions:write to the bot scopes in agent-team-manifest.json. NOTE: the new reactions:write scope requires Adam to re-apply the manifest to app A0BCC7TTU66 and reinstall the app. Until then reactions.add returns missing_scope, which the listener swallows (the reaction silently no-ops) — answer handling is unaffected. AUTHZ-01 and the first-answer-wins CAS remain unchanged. --- agent-team/slack/agent-team-manifest.json | 3 +- agent-team/tests/test_slack_listener.py | 191 ++++++++++++++++++++++ agent-team/tests/test_slack_live.py | 49 ++++++ 3 files changed, 242 insertions(+), 1 deletion(-) diff --git a/agent-team/slack/agent-team-manifest.json b/agent-team/slack/agent-team-manifest.json index 6345972..7b0ead9 100644 --- a/agent-team/slack/agent-team-manifest.json +++ b/agent-team/slack/agent-team-manifest.json @@ -26,7 +26,8 @@ "groups:history", "im:history", "app_mentions:read", - "commands" + "commands", + "reactions:write" ] } }, diff --git a/agent-team/tests/test_slack_listener.py b/agent-team/tests/test_slack_listener.py index 9b7cadc..f135c4a 100644 --- a/agent-team/tests/test_slack_listener.py +++ b/agent-team/tests/test_slack_listener.py @@ -687,3 +687,194 @@ def test_handle_event_swallows_sqlite_error_from_submit_answer( # Must not raise; returns None. assert listener.handle_event(_interactive_payload("q1")) is None + + +# --------------------------------------------------------------------------- +# 👍 reaction on received messages (after AUTHZ-01) — best-effort. +# --------------------------------------------------------------------------- + + +class _RecordingReactor: + """Records (channel, ts) reaction-add calls; optionally raises to test swallow.""" + + def __init__(self, *, boom: bool = False) -> None: + self.calls: list[tuple[str, str]] = [] + self._boom = boom + + def __call__(self, channel: str, ts: str) -> None: + self.calls.append((channel, ts)) + if self._boom: + raise RuntimeError("missing_scope: reactions:write not granted") + + +def _listener_with_reactor( + db_path: Path, + enqueue: Any, + reactor: Any, + *, + owner_ids: set[str] | None = frozenset({OWNER_ID}), +) -> SlackListener: + return SlackListener( + SlackTransport(channel="C123"), + db_path, + enqueue, + owner_ids=set(owner_ids) if owner_ids else None, + reactor=reactor, + ) + + +def test_reaction_added_to_thread_reply_after_authz(db_path: Path) -> None: + """An accepted thread-reply answer gets a 👍 on the reply message (event ts).""" + _seed_open_question(db_path, question_id="q1", channel_ref=SEED_CHANNEL_REF) + queue = RecordingQueue() + reactor = _RecordingReactor() + listener = _listener_with_reactor(db_path, queue, reactor) + + reply = _events_api_reply(text="approve", thread_ts=SEED_CHANNEL_REF) + outcome = listener.handle_event(reply) + + assert outcome is not None and outcome.accepted is True + # Reacted to the REPLY message: channel + the event ts (not the thread_ts). + assert reactor.calls == [("C123", "1700000001.000200")] + + +def test_reaction_not_added_for_non_owner(db_path: Path) -> None: + """AUTHZ-01 runs first: a non-owner message is rejected AND gets no reaction.""" + _seed_open_question(db_path, question_id="q1", channel_ref=SEED_CHANNEL_REF) + queue = RecordingQueue() + reactor = _RecordingReactor() + listener = _listener_with_reactor(db_path, queue, reactor) + + outcome = listener.handle_event( + _events_api_reply(sender_id="U_INTRUDER", thread_ts=SEED_CHANNEL_REF) + ) + + assert outcome is None + assert queue.jobs == [] + # The owner check returned before _maybe_react ran: NO reaction attempted. + assert reactor.calls == [] + assert _row_status(db_path, "q1") == "open" + + +def test_reaction_error_is_swallowed(db_path: Path) -> None: + """A reactor failure (e.g. missing reactions:write scope) never breaks handling.""" + _seed_open_question(db_path, question_id="q1", channel_ref=SEED_CHANNEL_REF) + queue = RecordingQueue() + reactor = _RecordingReactor(boom=True) + listener = _listener_with_reactor(db_path, queue, reactor) + + # Must not raise; the answer is still accepted + enqueued despite the reaction + # failing (the reaction silently no-ops until the scope is granted). + outcome = listener.handle_event( + _events_api_reply(text="approve", thread_ts=SEED_CHANNEL_REF) + ) + + assert outcome is not None and outcome.accepted is True + assert len(queue.jobs) == 1 + assert reactor.calls == [("C123", "1700000001.000200")] + + +def test_no_reactor_configured_is_quiet_noop(db_path: Path) -> None: + """With no reactor injected, an accepted answer simply does not react.""" + _seed_open_question(db_path, question_id="q1", channel_ref=SEED_CHANNEL_REF) + queue = RecordingQueue() + listener = _listener(db_path, queue, owner_ids={OWNER_ID}) # no reactor + + outcome = listener.handle_event( + _events_api_reply(text="approve", thread_ts=SEED_CHANNEL_REF) + ) + assert outcome is not None and outcome.accepted is True + + +def test_reaction_skipped_for_slash_or_interactive(db_path: Path) -> None: + """No reactable message ts on slash/interactive payloads => no reaction.""" + _seed_open_question(db_path, question_id="q1", thread_id="t1", turn=0) + queue = RecordingQueue() + reactor = _RecordingReactor() + listener = _listener_with_reactor(db_path, queue, reactor) + + # A block_actions interactive payload (no inner event ts) is accepted but has + # no reactable message, so no reaction is attempted. + outcome = listener.handle_event(_interactive_payload("q1", value="approve")) + assert outcome is not None and outcome.accepted is True + assert reactor.calls == [] + + +# --------------------------------------------------------------------------- +# /new-task — one-thread-per-task root "📥 Task received" ack post. +# --------------------------------------------------------------------------- + + +def test_new_task_posts_root_and_passes_root_ts(db_path: Path) -> None: + """A /new-task posts the root ack and forwards its ts to the callback. + + The "📥 Task received" message IS the slash command's acknowledgement (no + reaction), and its ``ts`` is threaded into start via the callback so every + later question/notification lands in one Slack thread. + """ + posted: list[dict[str, Any]] = [] + + def _poster(message: dict[str, Any]) -> dict[str, Any]: + posted.append(message) + return {"ts": "1700000000.ROOT"} + + seen: list[tuple[str, str, str]] = [] + + def _cb(task_text: str, via: str, root_ts: str) -> str: + seen.append((task_text, via, root_ts)) + return "thread-abc" + + listener = SlackListener( + SlackTransport(channel="C_TASK", poster=_poster), + db_path, + RecordingQueue(), + owner_ids={OWNER_ID}, + new_task_callback=_cb, + ) + + payload = { + "type": "slash_commands", + "command": "/new-task", + "text": "Add OAuth to the admin portal", + "user_id": OWNER_ID, + } + assert listener.handle_event(payload) is None + + # The root ack was posted to the channel with the 📥 prefix + description. + assert len(posted) == 1 + assert posted[0]["channel"] == "C_TASK" + assert posted[0]["text"].startswith("📥 Task received:") + assert "Add OAuth to the admin portal" in posted[0]["text"] + # The callback received the description, via, AND the captured root ts. + assert seen == [("Add OAuth to the admin portal", "slack", "1700000000.ROOT")] + + +def test_new_task_root_post_failure_degrades_to_empty_root_ts(db_path: Path) -> None: + """A failed root ack still starts the task (un-threaded): root_ts == ''.""" + + def _boom_poster(_message: dict[str, Any]) -> dict[str, Any]: + raise RuntimeError("slack down") + + seen: list[tuple[str, str, str]] = [] + + def _cb(task_text: str, via: str, root_ts: str) -> str: + seen.append((task_text, via, root_ts)) + return "thread-abc" + + listener = SlackListener( + SlackTransport(channel="C_TASK", poster=_boom_poster), + db_path, + RecordingQueue(), + owner_ids={OWNER_ID}, + new_task_callback=_cb, + ) + + payload = { + "type": "slash_commands", + "command": "/new-task", + "text": "do the thing", + "user_id": OWNER_ID, + } + assert listener.handle_event(payload) is None + # Task still started, with an empty root_ts (top-level questions). + assert seen == [("do the thing", "slack", "")] diff --git a/agent-team/tests/test_slack_live.py b/agent-team/tests/test_slack_live.py index 6d3fbf8..5911c89 100644 --- a/agent-team/tests/test_slack_live.py +++ b/agent-team/tests/test_slack_live.py @@ -38,11 +38,16 @@ class _FakeWebClient: response if response is not None else {"ts": "169.1", "ok": True} ) self.calls: list[dict[str, Any]] = [] + self.reaction_calls: list[dict[str, Any]] = [] def chat_postMessage(self, **kwargs: Any) -> dict[str, Any]: self.calls.append(kwargs) return self.response + def reactions_add(self, **kwargs: Any) -> dict[str, Any]: + self.reaction_calls.append(kwargs) + return {"ok": True} + class _DataResponse: """A ``slack_sdk.SlackResponse``-like object exposing the payload via ``.data``.""" @@ -161,6 +166,50 @@ def test_callback_id_dropped_metadata_carries_question_id() -> None: assert "blocks" in kwargs +def test_poster_forwards_thread_ts_to_chat_post_message() -> None: + """One-thread-per-task: thread_ts IS a postMessage param and is forwarded.""" + from agent_team.transport.slack_live import build_slack_poster as _bp + + client = _FakeWebClient() + transport = SlackTransport("C123", poster=_bp(client=client)) + + transport.post_question( + thread_id="task-7", + question_id="q-42", + turn=1, + question_set=_question_set(), + deadline="2026-06-18T00:00:00Z", + thread_ts="1700000000.ROOT", + ) + + assert client.calls[0]["thread_ts"] == "1700000000.ROOT" + + +def test_reactor_calls_reactions_add_with_thumbsup() -> None: + """build_slack_reactor performs reactions.add(channel, ts, name='thumbsup').""" + from agent_team.transport.slack_live import build_slack_reactor + + client = _FakeWebClient() + reactor = build_slack_reactor(client=client) + + reactor("C123", "1700000001.000200") + + assert client.reaction_calls == [ + {"channel": "C123", "timestamp": "1700000001.000200", "name": "thumbsup"} + ] + + +def test_reactor_missing_token_raises_runtime_error( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """No token + no client => the same loud RuntimeError as the poster path.""" + from agent_team.transport.slack_live import build_slack_reactor + + monkeypatch.delenv("SLACK_BOT_TOKEN", raising=False) + with pytest.raises(RuntimeError): + build_slack_reactor() + + def test_convenience_transport_factory() -> None: """``build_live_slack_transport`` wires the live poster onto a transport.""" client = _FakeWebClient()