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/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index 623899d..74ffc45 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], str] | None" = None, ) -> Any: """Build the live :class:`SlackListener` from the coordinator's seams (D-1). @@ -166,21 +167,47 @@ 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. 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, 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, + reactor=reactor, ) -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 @@ -199,7 +226,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, @@ -434,6 +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], str] | None" = None, + notify: "Callable[..., None] | None" = None, ) -> None: self._db_path = Path(db_path) self._transport = transport @@ -462,6 +497,15 @@ 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 + # 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 @@ -478,6 +522,22 @@ class Coordinator: # Accessors (the shared queue is the slack_listener handoff seam). # ------------------------------------------------------------------ # + def set_new_task_callback( + 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, 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 + @property def resume_queue(self) -> "queue.Queue[Any]": """The shared resume-job queue (slack_listener enqueues, drainer drains).""" @@ -582,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 @@ -591,17 +653,24 @@ 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. + + **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()") @@ -614,7 +683,12 @@ 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, + slack_thread_ts=slack_thread_ts, + ) question = graph_mod.pending_question(self._graph, thread_id=thread_id) if question is None: @@ -634,6 +708,7 @@ class Coordinator: self._transport, question_set, deadline=deadline, + thread_ts=slack_thread_ts or None, ) finally: conn.close() @@ -708,6 +783,164 @@ class Coordinator: # Maintenance tick (deadline policy + drain). # ------------------------------------------------------------------ # + 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) + + 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] + + # 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}`)' + + # 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 + 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. 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: + responder_mod.notify_question( + conn, + self._transport, + question["question_set"], + deadline=question.get("deadline") + or self._default_deadline(), + thread_ts=root_ts, + ) + 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"❓ {label} — needs more input. A new clarifying question was " + "posted above; reply in its thread.", + thread_ts=root_ts, + ) + continue + + # 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") + # 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( + 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.", + thread_ts=root_ts, + ) + else: + self._emit( + f"✅ {label} — plan ready for review (phase: {phase}).", + thread_ts=root_ts, + ) + + @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). @@ -727,7 +960,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). @@ -825,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() @@ -1028,6 +1272,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( @@ -1150,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 7c4b996..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, } ) @@ -591,6 +596,8 @@ 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). @@ -600,6 +607,20 @@ 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. + + ``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. @@ -610,6 +631,8 @@ def start_task( thread_id=tid, 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/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/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/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/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/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 2e09a91..7a26a03 100644 --- a/agent-team/agent_team/task_model.py +++ b/agent-team/agent_team/task_model.py @@ -87,6 +87,11 @@ class TaskRecord: thread_id: str status: TaskStatus 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) @@ -108,6 +113,16 @@ 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 + # 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] @@ -133,6 +148,8 @@ 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", ""), + 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 a35b305..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. @@ -523,6 +644,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() 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 7d07e3a..d699217 100644 --- a/agent-team/run-team.py +++ b/agent-team/run-team.py @@ -61,12 +61,13 @@ from __future__ import annotations import argparse import getpass import json +import logging import os 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 @@ -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, ...] = ( @@ -492,6 +494,72 @@ 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_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. @@ -507,20 +575,54 @@ 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 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() + + # 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. - return Coordinator( + coordinator = Coordinator( db_path=args.db, transport=transport, - build_plan_node=default_plan_node_factory, + build_clarify_node=lambda: default_clarify_node_factory( + context_provider=context_provider + ), + build_plan_node=lambda: default_plan_node_factory( + 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 + # 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, slack_thread_ts: coordinator.start_task( + task_text=task_text, + transport_name="slack", + slack_thread_ts=slack_thread_ts, + ) + ) + return coordinator def _build_transport(args: argparse.Namespace) -> Any: 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..aab77a6 --- /dev/null +++ b/agent-team/scripts/deploy-r720-ws-rollout.sh @@ -0,0 +1,109 @@ +#!/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 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. +# 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." diff --git a/agent-team/slack/agent-team-manifest.json b/agent-team/slack/agent-team-manifest.json index 64c1e30..7b0ead9 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,9 @@ "channels:history", "groups:history", "im:history", - "app_mentions:read" + "app_mentions:read", + "commands", + "reactions:write" ] } }, diff --git a/agent-team/tests/test_coordinator.py b/agent-team/tests/test_coordinator.py index 0692e35..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) @@ -912,3 +964,215 @@ 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, thread_ts=None: posted.append( + (qset, deadline, thread_ts) + ), + ) + + 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_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 + + 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 + + +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 + + +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_graph.py b/agent-team/tests/test_graph.py index 2a29dbc..b4b58be 100644 --- a/agent-team/tests/test_graph.py +++ b/agent-team/tests/test_graph.py @@ -136,6 +136,39 @@ 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_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_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 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_run_team.py b/agent-team/tests/test_run_team.py index 86b3574..2826bc5 100644 --- a/agent-team/tests/test_run_team.py +++ b/agent-team/tests/test_run_team.py @@ -644,22 +644,32 @@ class _FakeCoordinator: *, db_path: Any, transport: Any, + 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 self.build_plan_node = build_plan_node 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_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_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() 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 new file mode 100644 index 0000000..9e90760 --- /dev/null +++ b/agent-team/tests/test_ws_activation_wiring.py @@ -0,0 +1,229 @@ +"""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, 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 + # 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", + "slack_thread_ts": "1700000000.000100", + } + + +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, r: "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, r: "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 + + +# --------------------------------------------------------------------------- # +# 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" + ) 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,