From 253e31b0e85093b87b857ee0698428342b927aef Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Thu, 18 Jun 2026 12:56:42 -0400 Subject: [PATCH] feat(agent-team): live Slack transport + Socket Mode listener with owner allowlist slack_live: real slack_sdk poster. slack_listener: Socket Mode inbound; trust boundary = app-token auth + an explicit owner allowlist on the sender (fail-closed, rejects all if AGENT_TEAM_SLACK_OWNER_IDS unset) + the open-status CAS as anti-replay. Closes the AUTHZ-01 missing-sender-authz finding from the security review. --- .../agent_team/transport/slack_listener.py | 385 +++++++++++++++ agent-team/agent_team/transport/slack_live.py | 142 ++++++ agent-team/tests/test_slack_listener.py | 454 ++++++++++++++++++ agent-team/tests/test_slack_live.py | 257 ++++++++++ 4 files changed, 1238 insertions(+) create mode 100644 agent-team/agent_team/transport/slack_listener.py create mode 100644 agent-team/agent_team/transport/slack_live.py create mode 100644 agent-team/tests/test_slack_listener.py create mode 100644 agent-team/tests/test_slack_live.py diff --git a/agent-team/agent_team/transport/slack_listener.py b/agent-team/agent_team/transport/slack_listener.py new file mode 100644 index 0000000..b592a56 --- /dev/null +++ b/agent-team/agent_team/transport/slack_listener.py @@ -0,0 +1,385 @@ +"""Socket Mode inbound Slack listener for the human-gate responder (design §3.3.1). + +This is the inbound counterpart to :class:`~agent_team.transport.slack_adapter.SlackTransport`. +The adapter posts question-sets outbound; this listener receives Adam's answers +and drives them into the durable first-answer-wins compare-and-set: + + inbound Slack event + -> SlackTransport.parse_answer (normalize to question_id/answer/via) + -> submit_answer (atomic UPDATE ... WHERE status='open') + -> enqueue_resume(job) [only if accepted] + +The listener does NOT resume the LangGraph graph itself; its sole job is +normalize -> submit -> enqueue. The turn-guarded :class:`ResumeWorker` (owned by +the coordinator) consumes the enqueued :class:`~agent_team.responder.ResumeJob`. + +Transport choice — Socket Mode (NOT a public webhook). The R720 box is VPN-only, +so there is no public HTTPS endpoint to expose. Slack's Socket Mode opens an +*outbound* WebSocket from the box to Slack and authenticates with an app-level +token; events arrive over that authenticated socket. This is the only place the +network/SDK is touched, and the SDK import is deferred (``slack_bolt`` / +``slack_sdk`` are not installed in the Mac/test env), so this module imports +cleanly without them and :meth:`SlackListener.handle_event` is fully unit +testable with no socket. + +SECURITY + This module handles inbound UNTRUSTED Slack input plus auth. The trust + boundary is enforced by THREE independent layers, all of which must hold for + an inbound event to mutate the ledger: + + * (a) Socket Mode app-token authentication on the socket. Under Socket Mode + there is no inbound HTTP request, so there is no ``X-Slack-Signature`` to + verify; the transport itself is authenticated by the app-level token used + to open the outbound WebSocket (only a holder of that token can establish + the socket and receive events at all). + + * (b) An explicit owner allowlist on the SENDER (AUTHZ-01, CWE-862). The + trust model is single-owner (Adam): only an authorized Slack user id may + answer/steer the autonomous pipeline. ``handle_event`` extracts the + inbound sender's Slack user id and rejects the event (``return None``, + WITHOUT calling :func:`~agent_team.responder.submit_answer`) unless that id + is in the configured ``owner_ids`` allowlist. This FAILS CLOSED: if the + allowlist is empty / unconfigured, EVERY answer is rejected, and if no + sender id can be recovered the event is treated as unauthorized. Socket + membership alone is NOT authorization — any member of a channel the app is + in could otherwise win the first-answer-wins race. This layer is the fix + for the prior (insufficient) assumption that "maps to an open row" was + itself authorization. + + * (c) The open-status compare-and-set as ANTI-REPLAY (not authorization). + An authorized event's embedded ``question_id`` only has effect if it maps + to a real, still-``open`` ledger row, because + :func:`~agent_team.responder.submit_answer` runs + ``UPDATE ... WHERE question_id=? AND status='open'``. A replayed or stale + ``question_id`` for a closed / expired / superseded / nonexistent row + loses that compare-and-set (rowcount 0) and is a no-op + (``accepted=False``) — it can never resume a graph or overwrite an + existing answer. First-answer-wins also neutralizes duplicate redelivery. + This is anti-replay AFTER authorization, never a substitute for it. + + * Answers are DATA, never code. The answer value is extracted by the + transport and stored verbatim as JSON (``json.dumps`` in the responder). + This module never ``eval``s, executes, interpolates, or otherwise + interprets answer content — it is treated purely as opaque payload data. + + * Defensive event filtering. ``handle_event`` validates that the payload is + a mapping carrying a recoverable ``question_id`` before doing any work, and + swallows the :class:`ValueError` that :meth:`SlackTransport.parse_answer` + raises for an unrecoverable id. Unrelated / malformed events are + logged-and-ignored (``return None``) rather than crashing the listen loop, + so a hostile or noisy event stream cannot take the listener down. +""" + +from __future__ import annotations + +import logging +from collections.abc import Mapping +from pathlib import Path +from typing import Any + +from agent_team.db.schema import connect +from agent_team.responder import AnswerOutcome, EnqueueResume, submit_answer +from agent_team.transport.slack_adapter import SlackTransport + +__all__ = [ + "SlackListener", +] + +_LOG = logging.getLogger(__name__) + +# Inbound Slack event ``type`` values that can carry an answer for the human +# gate: an interactive Block Kit callback (button / select), a thread reply or +# mention message, or a slash-command invocation. Anything else (presence +# changes, channel joins, reactions, ...) is ignored. ``parse_answer`` does the +# real question_id recovery; this is a cheap first filter so unrelated events +# never reach it. +_ANSWER_BEARING_TYPES: frozenset[str] = frozenset( + { + "block_actions", + "message", + "app_mention", + "slash_commands", + "view_submission", + } +) + + +class SlackListener: + """Socket Mode inbound listener that drives answers into the responder. + + Constructed with the injected collaborators so it is fully unit-testable + with no SDK and no socket: + + * ``transport`` — the :class:`SlackTransport` whose ``parse_answer`` + normalizes an inbound payload to ``(question_id, answer, via)``; + * ``db_path`` — the agent-team SQLite file; a fresh connection is opened per + event (and closed) so the per-event compare-and-set is isolated; + * ``enqueue_resume`` — the coordinator's resume-queue ``put`` callable; an + accepted answer hands its :class:`~agent_team.responder.ResumeJob` to it; + * ``app_token`` / ``bot_token`` — optional injected Slack tokens used only by + :meth:`serve` to open the Socket Mode connection. Never required for + :meth:`handle_event`. + * ``owner_ids`` — the allowlist of authorized Slack user ids (the + single-owner trust model, AUTHZ-01). Only a sender whose id is in this set + may answer. If ``None``/empty the listener FAILS CLOSED and rejects every + answer; :meth:`serve` sources it from ``AGENT_TEAM_SLACK_OWNER_IDS`` when + not injected. + + The listener never resumes the graph; it only normalizes, submits, and + enqueues. See the module SECURITY note for the trust boundary. + """ + + def __init__( + self, + transport: SlackTransport, + db_path: Path | str, + enqueue_resume: EnqueueResume, + *, + app_token: str | None = None, + bot_token: str | None = None, + owner_ids: set[str] | None = None, + ) -> None: + self._transport = transport + self._db_path = Path(db_path) + self._enqueue_resume = enqueue_resume + self._app_token = app_token + self._bot_token = bot_token + # The owner allowlist (AUTHZ-01). An empty set is the fail-closed default: + # an unconfigured deploy rejects every answer. + self._owner_ids: set[str] = set(owner_ids) if owner_ids else set() + + def handle_event(self, raw_payload: Any) -> AnswerOutcome | None: + """Normalize + submit one inbound event; return its outcome or ``None``. + + Steps: + + 1. Defensively validate the payload is a mapping for an answer-bearing + event type that carries a recoverable ``question_id``. Unrelated or + malformed events are logged and ignored (``return None``) — never + raised — so the listen loop cannot be crashed by a hostile or noisy + event. + 1b. AUTHORIZE THE SENDER (AUTHZ-01, fail-closed). Extract the inbound + sender's Slack user id and require it to be in the configured owner + allowlist BEFORE any ledger work. If the allowlist is unconfigured, + or no sender id can be recovered, or the sender is not an owner, the + event is rejected (``return None``, ``submit_answer`` is NOT called). + 2. Open a per-event ledger connection, run + :func:`~agent_team.responder.submit_answer` (the atomic + first-answer-wins compare-and-set), then close the connection. + 3. If the answer was accepted (it was the first valid answer for a + still-``open`` ledger row), hand the resulting + :class:`~agent_team.responder.ResumeJob` to the injected + ``enqueue_resume`` via the responder; a duplicate / late / forged id + yields ``accepted=False`` and is a no-op. + + Returns the :class:`~agent_team.responder.AnswerOutcome` from + ``submit_answer``, or ``None`` if the event was not an answer we act on. + """ + if not isinstance(raw_payload, Mapping): + _LOG.debug("ignoring non-mapping Slack event: %r", type(raw_payload)) + return None + + event_type = raw_payload.get("type") + if event_type is not None and event_type not in _ANSWER_BEARING_TYPES: + _LOG.debug("ignoring non-answer Slack event type: %r", event_type) + return None + + # AUTHZ-01 (CWE-862), fail-closed: only an allowlisted owner may answer. + # Resolve the question_id first (best-effort) so a rejection log names the + # question without leaking answer content. parse_answer's id recovery is + # the same one submit_answer uses; an unrecoverable id is handled below. + if not self._is_authorized(raw_payload): + return None + + # ``submit_answer`` calls ``transport.parse_answer`` internally, which + # raises ValueError when no question_id is recoverable. Wrap the whole + # submit so a malformed / unrelated event is logged-and-ignored rather + # than crashing the loop. The enqueue itself happens inside + # ``submit_answer`` (only on accept), so it is covered by this guard too. + conn = connect(self._db_path) + try: + outcome = submit_answer( + conn, + self._transport, + raw_payload, + enqueue_resume=self._enqueue_resume, + ) + except ValueError as exc: + # Unrecoverable / forged-shape payload: parse_answer rejected it. + # This is expected for unrelated chatter on the channel; ignore it. + _LOG.debug("ignoring Slack event with no recoverable answer: %s", exc) + return None + finally: + conn.close() + + if outcome.accepted: + _LOG.info("accepted Slack answer for question_id=%s", outcome.question_id) + else: + _LOG.info( + "ignored Slack answer for question_id=%s (not open: duplicate / " + "late / forged id loses the compare-and-set)", + outcome.question_id, + ) + return outcome + + def _is_authorized(self, raw_payload: Mapping[str, Any]) -> bool: + """Return ``True`` iff the payload's sender is an allowlisted owner. + + Fails closed (AUTHZ-01, CWE-862): + + * If the allowlist is empty / unconfigured, reject EVERY answer and log a + warning naming ``AGENT_TEAM_SLACK_OWNER_IDS`` so the misprovisioning is + obvious. An unconfigured deploy must accept answers from no one. + * If no sender id can be recovered, treat the event as unauthorized. + * If the sender id is not in the allowlist, reject it. + + Never logs the answer content or any token — only the (best-effort) + question id and the offending sender id, which are non-secret routing + identifiers. + """ + if not self._owner_ids: + _LOG.warning( + "rejecting Slack answer: owner allowlist is unconfigured " + "(set AGENT_TEAM_SLACK_OWNER_IDS); the listener fails closed and " + "accepts answers from no one until it is provisioned" + ) + return False + + sender_id = _extract_sender_id(raw_payload) + if sender_id is None: + _LOG.warning( + "rejecting Slack answer: no recoverable sender id in payload " + "(treated as unauthorized)" + ) + return False + + if sender_id not in self._owner_ids: + # Log the rejection WITHOUT the answer content or any token. + _LOG.warning( + "rejecting Slack answer for question_id=%s: unauthorized sender %r " + "(not in owner allowlist)", + _safe_question_id(self._transport, raw_payload), + sender_id, + ) + return False + + return True + + def serve(self) -> None: # pragma: no cover - live socket, not unit-tested + """Open the Socket Mode connection and forward events to ``handle_event``. + + Lazily imports ``slack_bolt`` (deferred so this module imports cleanly + without the SDK, mirroring ``graph.build_sqlite_checkpointer``), wires a + handler that forwards every inbound event to :meth:`handle_event`, and + blocks on the Socket Mode handler. This is the ONLY method that touches + the network and is intentionally not unit-tested against a live socket; + :meth:`handle_event` carries all the testable logic. + + Also sources the owner allowlist from ``AGENT_TEAM_SLACK_OWNER_IDS`` + (comma-separated Slack user ids) when one was not injected, so the + production entry is allowlist-aware. The listener still FAILS CLOSED if + the env var is unset/empty — :meth:`handle_event` rejects every answer. + + Raises :class:`RuntimeError` if the SDK package or the required tokens + are missing, so a misconfigured deploy fails loudly rather than silently + never receiving answers. + """ + # Deferred import (mirrors the SDK import discipline): keep ``os`` out of + # the module's import-time surface so this stays cleanly importable. + if not self._owner_ids: + import os + + raw = os.environ.get("AGENT_TEAM_SLACK_OWNER_IDS", "") + self._owner_ids = {uid.strip() for uid in raw.split(",") if uid.strip()} + + if not self._app_token or not self._bot_token: + raise RuntimeError( + "SlackListener.serve requires both an app-level token " + "(xapp-, Socket Mode) and a bot token (xoxb-); inject them via " + "SlackListener(..., app_token=..., bot_token=...)." + ) + + try: + from slack_bolt import App + from slack_bolt.adapter.socket_mode import SocketModeHandler + except ImportError as exc: + raise RuntimeError( + "slack_bolt is unavailable; install 'slack-bolt' to run the " + "Socket Mode listener (SlackListener.serve). Tests exercise " + "handle_event directly with no SDK." + ) from exc + + app = App(token=self._bot_token) + + # Block Kit interactions, messages, mentions, and slash commands all + # funnel through the same normalize -> submit -> enqueue path. Slack Bolt + # dispatches by event family, so register the relevant ones and forward + # the raw body unchanged; handle_event does the filtering + parsing. + def _forward(body: Mapping[str, Any]) -> None: + self.handle_event(body) + + @app.action({}) # any block_actions interaction + def _on_action(ack: Any, body: Mapping[str, Any]) -> None: + ack() + _forward(body) + + @app.event("message") + def _on_message(body: Mapping[str, Any]) -> None: + _forward(body) + + @app.event("app_mention") + def _on_mention(body: Mapping[str, Any]) -> None: + _forward(body) + + SocketModeHandler(app, self._app_token).start() + + +def _extract_sender_id(raw_payload: Mapping[str, Any]) -> str | None: + """Recover the inbound sender's Slack user id from any supported shape. + + Handles the inbound payload shapes defensively (AUTHZ-01): + + * interactive ``block_actions`` / view submissions: ``payload["user"]["id"]``; + * Events API message / mention: ``payload["event"]["user"]`` (the event + envelope) or a top-level ``payload["user"]`` (which may itself be the id + string or a ``{"id": ...}`` mapping); + * slash command: ``payload["user_id"]``. + + Returns the non-empty user id string, or ``None`` if no sender id can be + recovered (the caller treats that as unauthorized — fail closed). + """ + # Interactive / top-level ``user``: a mapping ({"id": ...}) or a bare id str. + user = raw_payload.get("user") + if isinstance(user, Mapping): + uid = user.get("id") + if uid: + return str(uid) + elif isinstance(user, str) and user: + return user + + # Events API envelope: the inner event carries the author's user id. + event = raw_payload.get("event") + if isinstance(event, Mapping): + uid = event.get("user") + if uid: + return str(uid) + + # Slash command shape. + user_id = raw_payload.get("user_id") + if user_id: + return str(user_id) + + return None + + +def _safe_question_id(transport: SlackTransport, raw_payload: Mapping[str, Any]) -> str: + """Best-effort recover the question id for a rejection log, never raising. + + Used only to name the question in an unauthorized-sender warning. Returns a + placeholder if the id is unrecoverable; never logs answer content or tokens. + """ + try: + question_id, _answer, _via = transport.parse_answer(raw_payload) + except Exception: # noqa: BLE001 — logging path must never raise + return "" + return question_id diff --git a/agent-team/agent_team/transport/slack_live.py b/agent-team/agent_team/transport/slack_live.py new file mode 100644 index 0000000..9f9868e --- /dev/null +++ b/agent-team/agent_team/transport/slack_live.py @@ -0,0 +1,142 @@ +"""Live ``slack_sdk``-backed Slack poster (design §3.3.1, §7.1 P1 — Slack first). + +The :mod:`agent_team.transport.slack_adapter` module ships the §3.3.1 transport +contract with a dependency-injected ``poster`` seam: the adapter renders the +question-set into a message dict and hands it to a +``SlackPoster = Callable[[dict[str, Any]], Mapping[str, Any]]`` whose job is to +perform the real ``chat.postMessage`` and return a response carrying the message +``ts``. The foundation's default poster refuses the network so nothing ships +provisioned; this module supplies the **production** poster, backed by +``slack_sdk.WebClient``, that the P1 (Slack first) live wiring injects. + +Deferred import (mirrors :func:`agent_team.graph.build_sqlite_checkpointer`): +``slack_sdk`` is an optional dependency that may be absent in pre-deploy / +test environments, so this module imports cleanly without it. The import is +deferred to the moment a live client is actually constructed, and a missing +package raises a clear :class:`RuntimeError` so a misconfigured deploy fails +loudly rather than silently. + +Message-dict to ``chat.postMessage`` mapping +-------------------------------------------- +The adapter's message dict (see ``SlackTransport.post_question``) carries +``channel``, ``callback_id``, ``text``, ``blocks`` and ``metadata``. Slack's +``chat.postMessage`` Web API method does **not** accept a top-level +``callback_id`` keyword argument (``callback_id`` is a legacy attachment / +interactive-component field, not a message-post parameter), so passing it +through verbatim would raise a ``TypeError`` / Slack ``invalid_arguments``. + +The durable inbound key is therefore carried by ``metadata`` instead: the +adapter embeds ``question_id`` under ``metadata.event_payload.question_id``, and +``slack_adapter._extract_question_id`` reads exactly that path off an inbound +message. ``chat.postMessage`` *does* accept ``metadata``, so forwarding it +preserves the inbound mapping. The poster consequently **drops** ``callback_id`` +from the postMessage kwargs and forwards only the parameters the Web API +accepts (``channel``, ``text``, ``blocks``, ``metadata``), letting ``metadata`` +do the question-id round-trip the adapter relies on. +""" + +from __future__ import annotations + +import os +from collections.abc import Mapping +from typing import Any + +from agent_team.transport.slack_adapter import SlackPoster, SlackTransport + +__all__ = [ + "build_live_slack_transport", + "build_slack_poster", +] + +# 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") + + +def build_slack_poster(token: str | None = None, *, client: Any = None) -> SlackPoster: + """Build a live ``slack_sdk``-backed :data:`SlackPoster` (§3.3.1, P1). + + The returned callable accepts the adapter's rendered message dict, performs + a ``chat.postMessage``, and returns the response as a mapping carrying the + message ``ts`` so ``SlackTransport._extract_ts`` can record the + ``channel_ref``. + + ``client`` (optional) injects a pre-built Slack client for testability; any + object exposing ``chat_postMessage(**kwargs)`` works. When omitted, a + ``slack_sdk.WebClient`` is constructed lazily from ``token`` (falling back to + the ``SLACK_BOT_TOKEN`` environment variable). The ``slack_sdk`` import is + deferred so this module imports cleanly without the optional package; a + missing package or a missing token raises a clear :class:`RuntimeError`. + + The poster maps the adapter's message dict to the Web API's accepted + parameters: it forwards ``channel``, ``text``, ``blocks`` and ``metadata`` + and **drops** ``callback_id`` (not a ``chat.postMessage`` parameter — the + ``question_id`` round-trips via ``metadata.event_payload`` instead). See the + module docstring for the full rationale. + """ + if client is None: + client = _build_web_client(token) + + def _poster(message: dict[str, Any]) -> Mapping[str, Any]: + kwargs = {key: message[key] for key in _POST_MESSAGE_KEYS if key in message} + response = client.chat_postMessage(**kwargs) + return _as_mapping(response) + + return _poster + + +def build_live_slack_transport( + channel: str, token: str | None = None, *, client: Any = None +) -> SlackTransport: + """Build a :class:`SlackTransport` wired to a live ``slack_sdk`` poster. + + Convenience constructor for the P1 live coordinator: equivalent to + ``SlackTransport(channel, poster=build_slack_poster(token, client=client))``. + See :func:`build_slack_poster` for the token / client / deferred-import + semantics. + """ + return SlackTransport(channel, poster=build_slack_poster(token, client=client)) + + +def _build_web_client(token: str | None) -> Any: + """Lazily construct a ``slack_sdk.WebClient`` (deferred optional import). + + Raises a clear :class:`RuntimeError` if ``slack_sdk`` is not installed or no + token is resolvable (neither ``token`` nor ``SLACK_BOT_TOKEN``), so a + misconfigured deploy fails loudly rather than silently. + """ + try: + from slack_sdk import WebClient + except ImportError as exc: # pragma: no cover - depends on optional dep + raise RuntimeError( + "slack_sdk is unavailable; install the 'slack_sdk' package to build " + "a live Slack poster (P1), or inject a 'client' for testing." + ) from exc + + resolved = token or os.environ.get("SLACK_BOT_TOKEN") + if not resolved: + raise RuntimeError( + "No Slack bot token available; pass 'token' or set the " + "SLACK_BOT_TOKEN environment variable to build a live Slack poster." + ) + return WebClient(token=resolved) + + +def _as_mapping(response: Any) -> Mapping[str, Any]: + """Coerce a ``chat_postMessage`` response to a plain mapping. + + ``slack_sdk`` returns a ``SlackResponse`` exposing the payload via ``.data``; + if a test injects a client returning a bare mapping, accept it as-is. The + result must carry ``ts`` so ``SlackTransport._extract_ts`` recovers the + ``channel_ref``. + """ + if isinstance(response, Mapping): + return response + data = getattr(response, "data", None) + if isinstance(data, Mapping): + return data + raise TypeError( + "Slack chat_postMessage returned an unsupported response; expected a " + f"mapping or an object with a mapping '.data', got {type(response)!r}" + ) diff --git a/agent-team/tests/test_slack_listener.py b/agent-team/tests/test_slack_listener.py new file mode 100644 index 0000000..10eaf85 --- /dev/null +++ b/agent-team/tests/test_slack_listener.py @@ -0,0 +1,454 @@ +"""Unit tests for agent_team.transport.slack_listener (§3.3.1). + +All mocked — no network, no slack_sdk / slack_bolt. Covers: + +* the module imports cleanly without the Slack SDK installed; +* handle_event on a valid interactive payload from an AUTHORIZED owner accepts + and enqueues a resume job; +* AUTHZ-01: a non-owner sender is rejected (submit_answer + enqueue NOT called, + ledger row stays open); an unconfigured allowlist rejects every answer + (fail-closed); a payload with no recoverable sender id is rejected; +* a duplicate event for the same question_id is a no-op (first-answer-wins) and + does NOT enqueue again; +* an unrelated / malformed event (no recoverable question_id) is ignored without + crashing and never enqueues; +* a forged question_id for a nonexistent / closed row is a no-op (accepted=False) + — the anti-replay layer (the responder's WHERE status='open' CAS). + +The accept / duplicate / forged cases drive a REAL on-disk SQLite ledger (the +foundation ``init_db`` / ``connect`` + a seeded open question via +``notify_question``) so the actual compare-and-set runs. +""" + +from __future__ import annotations + +import importlib +from pathlib import Path +from typing import Any + +import pytest + +from agent_team.db.schema import connect, init_db +from agent_team.responder import ResumeJob, notify_question +from agent_team.transport.base import QuestionSet, Transport +from agent_team.transport.slack_adapter import SlackTransport, build_callback_id +from agent_team.transport.slack_listener import SlackListener + + +# --------------------------------------------------------------------------- +# Fixtures + helpers. +# --------------------------------------------------------------------------- + +# The single authorized owner id used across the accept-path tests (AUTHZ-01). +OWNER_ID = "U_OWNER" + + +@pytest.fixture +def db_path(tmp_path: Path) -> Path: + """An initialized on-disk ledger DB file (foundation schema).""" + path = tmp_path / "agent-team.db" + init_db(path) + return path + + +class RecordingQueue: + """Captures the resume jobs the listener enqueues.""" + + def __init__(self) -> None: + self.jobs: list[ResumeJob] = [] + + def __call__(self, job: ResumeJob) -> None: + self.jobs.append(job) + + +def _seed_open_question( + db_path: Path, + *, + thread_id: str = "t1", + question_id: str = "q1", + turn: int = 0, +) -> None: + """Insert a real ``open`` ledger row via notify_question (no network post).""" + conn = connect(db_path) + try: + # An injected poster that returns a ts ref; never reaches the network. + transport = SlackTransport( + channel="C123", poster=lambda _msg: {"ts": "1700000000.000100"} + ) + qs = QuestionSet( + thread_id=thread_id, + question_id=question_id, + turn=turn, + questions=["proceed?"], + context={"repo": "x"}, + ) + notify_question(conn, transport, qs, deadline="2026-06-18T00:00:00+00:00") + finally: + conn.close() + + +def _interactive_payload( + question_id: str, + value: str = "approve", + *, + sender_id: str | None = OWNER_ID, +) -> dict[str, Any]: + """A minimal Slack ``block_actions`` payload carrying ``question_id``. + + Carries the interactive sender id (``user.id``) so the AUTHZ-01 allowlist + check can run. Pass ``sender_id=None`` to omit the sender entirely (the + no-recoverable-sender case). + """ + payload: dict[str, Any] = { + "type": "block_actions", + "callback_id": build_callback_id(question_id), + "actions": [{"action_id": "answer", "value": value}], + } + if sender_id is not None: + payload["user"] = {"id": sender_id} + return payload + + +def _listener( + db_path: Path, + enqueue: Any, + *, + owner_ids: set[str] | None = frozenset({OWNER_ID}), +) -> SlackListener: + """Construct a listener, authorized for ``OWNER_ID`` by default. + + Pass ``owner_ids=None`` (or an empty set) to exercise the fail-closed + unconfigured-allowlist path. + """ + return SlackListener( + SlackTransport(channel="C123"), + db_path, + enqueue, + owner_ids=set(owner_ids) if owner_ids else None, + ) + + +def _row_status(db_path: Path, question_id: str) -> str | None: + """Return the ledger ``status`` for ``question_id`` (or ``None`` if absent).""" + conn = connect(db_path) + try: + row = conn.execute( + "SELECT status FROM pending_questions WHERE question_id=?", + (question_id,), + ).fetchone() + finally: + conn.close() + return None if row is None else str(row["status"]) + + +# --------------------------------------------------------------------------- +# Import cleanliness (no SDK). +# --------------------------------------------------------------------------- + + +def test_module_imports_without_slack_sdk(monkeypatch: pytest.MonkeyPatch) -> None: + """Reloading the module + constructing a listener never imports the SDK. + + Mirrors ``tests/test_slack_live.py``: monkeypatch ``builtins.__import__`` to + raise ImportError for any ``slack_sdk`` / ``slack_bolt`` import, then reload + the module under test. This proves the SDK import is genuinely deferred (it + is touched only in :meth:`SlackListener.serve`, never at module import or + listener construction time) — independent of any prior ``sys.modules`` + state, unlike a global-state precondition that is merely order-dependent. + """ + import builtins + + import agent_team.transport.slack_listener as slack_listener_module + + real_import = builtins.__import__ + + def _blocked_import(name: str, *args: Any, **kwargs: Any) -> Any: + if ( + name == "slack_sdk" + or name.startswith("slack_sdk.") + or name == "slack_bolt" + or name.startswith("slack_bolt.") + ): + raise ImportError(f"{name} is blocked for this test") + return real_import(name, *args, **kwargs) + + monkeypatch.setattr(builtins, "__import__", _blocked_import) + + module = importlib.reload(slack_listener_module) + + # The reload succeeded with the SDK blocked, and the listener is + # constructible without ever importing slack_sdk / slack_bolt. + listener = module.SlackListener( + SlackTransport(channel="C123"), Path(":memory:"), lambda _job: None + ) + assert isinstance(listener, module.SlackListener) + + +# --------------------------------------------------------------------------- +# handle_event — accept + enqueue (real CAS). +# --------------------------------------------------------------------------- + + +def test_handle_event_accepts_and_enqueues(db_path: Path) -> None: + _seed_open_question(db_path, question_id="q1", thread_id="t1", turn=0) + queue = RecordingQueue() + listener = _listener(db_path, queue) + + outcome = listener.handle_event(_interactive_payload("q1", value="approve")) + + assert outcome is not None + assert outcome.accepted is True + assert outcome.question_id == "q1" + assert len(queue.jobs) == 1 + job = queue.jobs[0] + assert job.thread_id == "t1" + assert job.question_id == "q1" + assert job.turn == 0 + assert job.answer == "approve" + + +# --------------------------------------------------------------------------- +# handle_event — AUTHZ-01 owner allowlist (CWE-862), fail-closed. +# --------------------------------------------------------------------------- + + +def test_handle_event_rejects_non_owner_sender(db_path: Path) -> None: + """A sender not in the owner allowlist is rejected; the open row is untouched. + + The payload carries a valid, recoverable question_id mapping to a real open + ledger row, so ONLY the sender-identity check stands between the attacker and + the first-answer-wins CAS. Prove submit_answer is never reached: no enqueue, + and the seeded row stays ``open`` in the real sqlite ledger. + """ + _seed_open_question(db_path, question_id="q1", thread_id="t1", turn=0) + queue = RecordingQueue() + listener = _listener(db_path, queue, owner_ids={OWNER_ID}) + + payload = _interactive_payload("q1", value="approve", sender_id="U_INTRUDER") + outcome = listener.handle_event(payload) + + assert outcome is None + assert queue.jobs == [] + # The CAS never ran: the row is still open (submit_answer was not called). + assert _row_status(db_path, "q1") == "open" + + +def test_handle_event_rejects_when_allowlist_unconfigured(db_path: Path) -> None: + """Fail-closed: with no owner allowlist, EVERY answer is rejected. + + Even a valid question_id from an otherwise-plausible sender is rejected so an + unprovisioned deploy accepts answers from no one. + """ + _seed_open_question(db_path, question_id="q1") + queue = RecordingQueue() + listener = _listener(db_path, queue, owner_ids=None) # unconfigured + + outcome = listener.handle_event( + _interactive_payload("q1", value="approve", sender_id=OWNER_ID) + ) + + assert outcome is None + assert queue.jobs == [] + assert _row_status(db_path, "q1") == "open" + + +def test_handle_event_rejects_empty_allowlist(db_path: Path) -> None: + """An explicitly empty allowlist is also fail-closed (rejects everything).""" + _seed_open_question(db_path, question_id="q1") + queue = RecordingQueue() + listener = _listener(db_path, queue, owner_ids=set()) + + outcome = listener.handle_event(_interactive_payload("q1", sender_id=OWNER_ID)) + + assert outcome is None + assert queue.jobs == [] + assert _row_status(db_path, "q1") == "open" + + +def test_handle_event_rejects_payload_without_sender_id(db_path: Path) -> None: + """A payload from which no sender id can be recovered is rejected (unauthorized).""" + _seed_open_question(db_path, question_id="q1") + queue = RecordingQueue() + listener = _listener(db_path, queue, owner_ids={OWNER_ID}) + + # sender_id=None omits ``user`` entirely; no event/user_id either. + outcome = listener.handle_event(_interactive_payload("q1", sender_id=None)) + + assert outcome is None + assert queue.jobs == [] + assert _row_status(db_path, "q1") == "open" + + +def test_handle_event_accepts_events_api_owner_sender(db_path: Path) -> None: + """An Events API message shape resolves the sender via ``event.user``.""" + _seed_open_question(db_path, question_id="q1", thread_id="t1", turn=0) + queue = RecordingQueue() + listener = _listener(db_path, queue, owner_ids={OWNER_ID}) + + payload = { + "type": "message", + "callback_id": build_callback_id("q1"), + "answer": "approve", + "event": {"user": OWNER_ID, "type": "message"}, + } + outcome = listener.handle_event(payload) + + assert outcome is not None and outcome.accepted is True + assert len(queue.jobs) == 1 + + +def test_handle_event_rejects_events_api_non_owner(db_path: Path) -> None: + """An Events API message from a non-owner ``event.user`` is rejected.""" + _seed_open_question(db_path, question_id="q1") + queue = RecordingQueue() + listener = _listener(db_path, queue, owner_ids={OWNER_ID}) + + payload = { + "type": "message", + "callback_id": build_callback_id("q1"), + "answer": "approve", + "event": {"user": "U_INTRUDER", "type": "message"}, + } + outcome = listener.handle_event(payload) + + assert outcome is None + assert queue.jobs == [] + assert _row_status(db_path, "q1") == "open" + + +def test_handle_event_accepts_slash_command_owner_sender(db_path: Path) -> None: + """A slash-command shape resolves the sender via ``user_id``.""" + _seed_open_question(db_path, question_id="q1", thread_id="t1", turn=0) + queue = RecordingQueue() + listener = _listener(db_path, queue, owner_ids={OWNER_ID}) + + payload = { + "type": "slash_commands", + "callback_id": build_callback_id("q1"), + "text": "approve", + "user_id": OWNER_ID, + } + outcome = listener.handle_event(payload) + + assert outcome is not None and outcome.accepted is True + assert len(queue.jobs) == 1 + + +# --------------------------------------------------------------------------- +# handle_event — duplicate event (first-answer-wins). +# --------------------------------------------------------------------------- + + +def test_handle_event_duplicate_is_noop(db_path: Path) -> None: + """Second answer for the same question_id loses the CAS; no re-enqueue.""" + _seed_open_question(db_path, question_id="q1") + queue = RecordingQueue() + listener = _listener(db_path, queue) + + first = listener.handle_event(_interactive_payload("q1", value="approve")) + second = listener.handle_event(_interactive_payload("q1", value="reject")) + + assert first is not None and first.accepted is True + assert second is not None and second.accepted is False + # First-answer-wins: only the first answer enqueued a resume job. + assert len(queue.jobs) == 1 + assert queue.jobs[0].answer == "approve" + + +# --------------------------------------------------------------------------- +# handle_event — unrelated / malformed event is ignored. +# --------------------------------------------------------------------------- + + +def test_handle_event_ignores_non_mapping(db_path: Path) -> None: + queue = RecordingQueue() + listener = _listener(db_path, queue) + + assert listener.handle_event("not a mapping") is None + assert listener.handle_event(None) is None + assert queue.jobs == [] + + +def test_handle_event_ignores_unrelated_event_type(db_path: Path) -> None: + """An event of a type we never act on is filtered before parsing.""" + queue = RecordingQueue() + listener = _listener(db_path, queue) + + # A reaction event carries no question_id and is not answer-bearing. + outcome = listener.handle_event({"type": "reaction_added", "reaction": "thumbsup"}) + + assert outcome is None + assert queue.jobs == [] + + +def test_handle_event_ignores_payload_without_question_id(db_path: Path) -> None: + """An answer-bearing type with no recoverable question_id is logged-ignored.""" + queue = RecordingQueue() + listener = _listener(db_path, queue) + + # A message with no callback_id / metadata / question_id: parse_answer raises + # ValueError, which handle_event swallows. + outcome = listener.handle_event( + {"type": "message", "text": "just chatting", "channel": "C123"} + ) + + assert outcome is None + assert queue.jobs == [] + + +# --------------------------------------------------------------------------- +# handle_event — forged question_id (trust boundary held by the CAS). +# --------------------------------------------------------------------------- + + +def test_handle_event_forged_question_id_is_noop(db_path: Path) -> None: + """A well-formed payload whose question_id matches no open row is a no-op. + + The id is recoverable (so parse_answer succeeds), but it maps to no + ``open`` ledger row, so the responder's ``WHERE status='open'`` + compare-and-set returns rowcount 0 → accepted=False. This is the documented + trust boundary: a forged / replayed id cannot resume a graph. + """ + # Note: NO row seeded for this id. + queue = RecordingQueue() + listener = _listener(db_path, queue) + + outcome = listener.handle_event(_interactive_payload("forged-qid")) + + assert outcome is not None + assert outcome.accepted is False + assert outcome.question_id == "forged-qid" + assert queue.jobs == [] + + +def test_handle_event_closed_row_is_noop(db_path: Path) -> None: + """A second submit after the row is already answered loses the CAS too.""" + _seed_open_question(db_path, question_id="q1") + queue = RecordingQueue() + listener = _listener(db_path, queue) + + listener.handle_event(_interactive_payload("q1")) # closes the row + queue.jobs.clear() + + # Row is now 'answered'; a fresh forged event for it is a no-op. + outcome = listener.handle_event(_interactive_payload("q1", value="late")) + assert outcome is not None + assert outcome.accepted is False + assert queue.jobs == [] + + +# --------------------------------------------------------------------------- +# serve — token guard (no live socket). +# --------------------------------------------------------------------------- + + +def test_serve_requires_tokens(db_path: Path) -> None: + """serve raises a clear RuntimeError when tokens are missing.""" + listener = _listener(db_path, RecordingQueue()) + with pytest.raises(RuntimeError, match="app-level token"): + listener.serve() + + +def test_real_slack_transport_is_a_transport() -> None: + """Sanity: the injected SlackTransport is the contract the listener expects.""" + assert isinstance(SlackTransport(channel="C123"), Transport) diff --git a/agent-team/tests/test_slack_live.py b/agent-team/tests/test_slack_live.py new file mode 100644 index 0000000..48dfdf1 --- /dev/null +++ b/agent-team/tests/test_slack_live.py @@ -0,0 +1,257 @@ +"""Unit tests for agent_team.transport.slack_live (§3.3.1, §7.1 P1). + +The live poster is the production ``slack_sdk`` backing for the §3.3.1 injected +``SlackPoster`` seam. These tests prove the contract entirely with mocks (no +network, and ``slack_sdk`` itself is never required): the poster maps the +adapter's message dict to the ``chat.postMessage`` parameters Slack accepts, +the ``ts`` round-trips as the ``channel_ref`` through a real ``SlackTransport``, +``callback_id`` is dropped (not a valid postMessage param) while ``metadata`` +carries the durable ``question_id``, and a missing package / token or a client +failure fails loudly. +""" + +from __future__ import annotations + +import importlib +from typing import Any + +import pytest + +from agent_team.transport.base import QuestionSet +from agent_team.transport.slack_adapter import SlackPostError, SlackTransport +from agent_team.transport.slack_live import ( + build_live_slack_transport, + build_slack_poster, +) + + +# --------------------------------------------------------------------------- # +# Test doubles # +# --------------------------------------------------------------------------- # + + +class _FakeWebClient: + """A fake ``slack_sdk.WebClient`` recording ``chat_postMessage`` kwargs.""" + + def __init__(self, response: dict[str, Any] | None = None) -> None: + self.response = ( + response if response is not None else {"ts": "169.1", "ok": True} + ) + self.calls: list[dict[str, Any]] = [] + + def chat_postMessage(self, **kwargs: Any) -> dict[str, Any]: + self.calls.append(kwargs) + return self.response + + +class _DataResponse: + """A ``slack_sdk.SlackResponse``-like object exposing the payload via ``.data``.""" + + def __init__(self, data: dict[str, Any]) -> None: + self.data = data + + +class _SlackApiErrorLike(Exception): + """Stands in for ``slack_sdk.errors.SlackApiError`` (no slack_sdk needed).""" + + +class _FailingClient: + """A fake client whose ``chat_postMessage`` raises a Slack-API-like error.""" + + def chat_postMessage(self, **kwargs: Any) -> dict[str, Any]: + raise _SlackApiErrorLike("the_dog_ate_it") + + +def _question_set() -> QuestionSet: + return QuestionSet( + thread_id="task-7", + question_id="q-42", + turn=1, + questions=["Ship it?"], + context={"repo": "agent-team"}, + ) + + +# --------------------------------------------------------------------------- # +# Clean import without slack_sdk # +# --------------------------------------------------------------------------- # + + +def test_module_imports_without_slack_sdk(monkeypatch: pytest.MonkeyPatch) -> None: + """The module imports cleanly even when ``slack_sdk`` cannot be imported.""" + import builtins + + real_import = builtins.__import__ + + def _blocked_import(name: str, *args: Any, **kwargs: Any) -> Any: + if name == "slack_sdk" or name.startswith("slack_sdk."): + raise ImportError("slack_sdk is blocked for this test") + return real_import(name, *args, **kwargs) + + monkeypatch.setattr(builtins, "__import__", _blocked_import) + + module = importlib.reload( + importlib.import_module("agent_team.transport.slack_live") + ) + assert hasattr(module, "build_slack_poster") + assert hasattr(module, "build_live_slack_transport") + + +# --------------------------------------------------------------------------- # +# Happy path: injected fake client # +# --------------------------------------------------------------------------- # + + +def test_poster_returns_mapping_with_ts() -> None: + """The poster forwards to the client and returns a mapping carrying ``ts``.""" + client = _FakeWebClient() + poster = build_slack_poster(client=client) + + result = poster( + { + "channel": "C123", + "callback_id": "shq:q-42", + "text": "hi", + "blocks": [], + "metadata": {"event_type": "agent_team_question"}, + } + ) + + assert result["ts"] == "169.1" + + +def test_post_question_round_trips_ts_as_channel_ref() -> None: + """Wired through a real ``SlackTransport``, ``ts`` becomes the channel_ref.""" + client = _FakeWebClient() + transport = SlackTransport("C123", poster=build_slack_poster(client=client)) + + channel_ref = transport.post_question( + thread_id="task-7", + question_id="q-42", + turn=1, + question_set=_question_set(), + deadline="2026-06-18T00:00:00Z", + ) + + assert channel_ref == "169.1" + + +def test_callback_id_dropped_metadata_carries_question_id() -> None: + """``callback_id`` is not sent; ``metadata`` carries the durable question_id.""" + client = _FakeWebClient() + transport = SlackTransport("C123", poster=build_slack_poster(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", + ) + + assert len(client.calls) == 1 + kwargs = client.calls[0] + # callback_id is NOT a valid chat.postMessage parameter and must be dropped. + assert "callback_id" not in kwargs + # metadata IS forwarded and carries the durable inbound question_id. + assert kwargs["metadata"]["event_payload"]["question_id"] == "q-42" + # the accepted parameters are forwarded. + assert kwargs["channel"] == "C123" + assert "text" in kwargs + assert "blocks" in kwargs + + +def test_convenience_transport_factory() -> None: + """``build_live_slack_transport`` wires the live poster onto a transport.""" + client = _FakeWebClient() + transport = build_live_slack_transport("C123", client=client) + + channel_ref = transport.post_question( + thread_id="task-7", + question_id="q-42", + turn=1, + question_set=_question_set(), + deadline="2026-06-18T00:00:00Z", + ) + + assert channel_ref == "169.1" + + +def test_poster_accepts_slack_response_with_data_attr() -> None: + """A ``SlackResponse``-like object is coerced via its ``.data`` mapping.""" + client = _FakeWebClient(response=None) + client.response = _DataResponse({"ts": "169.1", "ok": True}) # type: ignore[assignment] + poster = build_slack_poster(client=client) + + result = poster({"channel": "C123", "text": "hi"}) + + assert result["ts"] == "169.1" + + +class _UnsupportedResponseClient: + """A fake client returning neither a mapping nor an object with ``.data``.""" + + def chat_postMessage(self, **kwargs: Any) -> Any: + return object() + + +def test_poster_unsupported_response_raises_type_error() -> None: + """A response that is neither a mapping nor has a mapping ``.data`` is fatal. + + Exercises the ``_as_mapping`` guard: a bare object (no ``ts``, no ``.data``) + cannot yield a ``channel_ref``, so the poster raises ``TypeError`` with the + guard's "unsupported response" message rather than silently dropping the ts. + """ + poster = build_slack_poster(client=_UnsupportedResponseClient()) + + with pytest.raises(TypeError, match="unsupported response"): + poster({"channel": "C123", "text": "hi"}) + + +# --------------------------------------------------------------------------- # +# Failure modes # +# --------------------------------------------------------------------------- # + + +def test_missing_token_raises_runtime_error(monkeypatch: pytest.MonkeyPatch) -> None: + """No token and no SLACK_BOT_TOKEN raises a clear RuntimeError.""" + monkeypatch.delenv("SLACK_BOT_TOKEN", raising=False) + + with pytest.raises(RuntimeError, match="Slack bot token"): + build_slack_poster() + + +def test_missing_package_raises_runtime_error( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """A missing ``slack_sdk`` package raises a clear RuntimeError.""" + import builtins + + real_import = builtins.__import__ + + def _blocked_import(name: str, *args: Any, **kwargs: Any) -> Any: + if name == "slack_sdk" or name.startswith("slack_sdk."): + raise ImportError("slack_sdk is blocked for this test") + return real_import(name, *args, **kwargs) + + monkeypatch.setattr(builtins, "__import__", _blocked_import) + monkeypatch.setenv("SLACK_BOT_TOKEN", "xoxb-present") + + with pytest.raises(RuntimeError, match="slack_sdk is unavailable"): + build_slack_poster() + + +def test_client_failure_surfaces_as_slack_post_error() -> None: + """A SlackApiError-like client failure surfaces as ``SlackPostError``.""" + transport = SlackTransport( + "C123", poster=build_slack_poster(client=_FailingClient()) + ) + + with pytest.raises(SlackPostError): + transport.post_question( + thread_id="task-7", + question_id="q-42", + turn=1, + question_set=_question_set(), + deadline="2026-06-18T00:00:00Z", + )