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.
This commit is contained in:
Adam Moussa 2026-06-18 12:56:42 -04:00
parent 4b17e8ebd4
commit 253e31b0e8
4 changed files with 1238 additions and 0 deletions

View file

@ -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 "<unknown>"
return question_id

View file

@ -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}"
)

View file

@ -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)

View file

@ -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",
)