fix(agent-team): handle real slack_bolt event envelope + map thread replies; bind start invoker #26

Merged
amoussa1229 merged 2 commits from fix/agent-team-slack-inbound-envelope into main 2026-06-22 19:30:48 +00:00
4 changed files with 339 additions and 44 deletions

View file

@ -37,6 +37,7 @@ __all__ = [
"answer_question",
"connect",
"expire_question",
"find_open_question_by_channel_ref",
"init_db",
"migrate",
"reopen_question",
@ -297,6 +298,37 @@ def reopen_question(
)
def find_open_question_by_channel_ref(
conn: sqlite3.Connection,
channel_ref: str,
) -> str | None:
"""Map a transport ``channel_ref`` to its still-``open`` ``question_id``.
The clarifier posts a question message and stores that message's ``ts`` as
the ledger row's ``channel_ref``; an inbound thread reply carries that same
value as its ``thread_ts``. When a reply carries no explicit
``callback_id`` / ``question_id`` / metadata (the real free-text-reply
shape), this resolves WHICH question the reply answers by its thread anchor.
The lookup is CONSTRAINED to ``status='open'`` (anti-replay): a stale or
replayed ``thread_ts`` pointing at a closed / expired / answered /
superseded row resolves to ``None`` and is a no-op for the caller. This
only resolves which question a reply targets; it is NEVER authorization —
the caller authorizes the sender first and fails closed.
Returns the ``question_id`` of the matching open row, or ``None`` if
``channel_ref`` is empty or matches no open row.
"""
if not channel_ref:
return None
row = conn.execute(
"SELECT question_id FROM pending_questions "
"WHERE channel_ref=? AND status='open'",
(channel_ref,),
).fetchone()
return None if row is None else str(row["question_id"])
def supersede_question(
conn: sqlite3.Connection,
*,

View file

@ -78,7 +78,7 @@ from collections.abc import Mapping
from pathlib import Path
from typing import Any
from agent_team.db.schema import connect
from agent_team.db.schema import connect, find_open_question_by_channel_ref
from agent_team.responder import AnswerOutcome, EnqueueResume, submit_answer
from agent_team.transport.slack_adapter import SlackTransport
@ -88,12 +88,22 @@ __all__ = [
_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.
# Discriminating 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.
#
# IMPORTANT — real slack_bolt envelope shapes (the bug this gate now handles):
# * Events API (message / app_mention) arrives wrapped as
# ``{"type": "event_callback", "event": {"type": "message", ...}}`` — the
# DISCRIMINATING type is the INNER ``event.type`` (``message`` /
# ``app_mention``), not the outer ``event_callback``.
# * Interactive (block_actions / view_submission) and slash_commands carry
# their discriminating type at the TOP level.
# :func:`_discriminating_type` collapses both to the inner type so this gate and
# the downstream extraction operate on a consistent view. ``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",
@ -104,6 +114,10 @@ _ANSWER_BEARING_TYPES: frozenset[str] = frozenset(
}
)
# Outer Slack envelope ``type`` for an Events API delivery; the real event (and
# its discriminating type) is nested under ``event``.
_EVENT_CALLBACK_TYPE = "event_callback"
class SlackListener:
"""Socket Mode inbound listener that drives answers into the responder.
@ -184,29 +198,42 @@ class SlackListener:
_LOG.debug("ignoring non-mapping Slack event: %r", type(raw_payload))
return None
event_type = raw_payload.get("type")
# Collapse the real slack_bolt Events API envelope (``event_callback``
# wrapping an inner ``event``) to its discriminating inner type so the
# gate sees ``message`` / ``app_mention`` rather than ``event_callback``.
event_type = _discriminating_type(raw_payload)
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.
# This runs FIRST, before any question resolution or ledger work; it
# gates on the SENDER's Slack user id (``event.user`` for the Events API
# shape, ``user.id`` interactive, ``user_id`` slash). The thread-reply
# mapping below only resolves WHICH question is answered, never WHO may
# answer.
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.
# raises ValueError when no question_id is recoverable. A real free-text
# thread reply carries no callback_id / question_id / metadata, so its
# question is resolved here by ``thread_ts`` (== the question message's
# ``ts``, stored as the row's ``channel_ref``) and fed to the responder
# as an explicit normalized payload. Explicit id recovery takes
# precedence; the thread_ts fallback is only tried when parse_answer
# cannot recover an id. Wrap the whole submit so a malformed / unrelated
# event is logged-and-ignored rather than crashing the loop. The enqueue
# happens inside ``submit_answer`` (only on accept), so it is covered too.
conn = connect(self._db_path)
try:
payload = self._resolve_payload(conn, raw_payload)
if payload is None:
return None
outcome = submit_answer(
conn,
self._transport,
raw_payload,
payload,
enqueue_resume=self._enqueue_resume,
)
except ValueError as exc:
@ -270,6 +297,67 @@ class SlackListener:
return True
def _resolve_payload(
self, conn: Any, raw_payload: Mapping[str, Any]
) -> Mapping[str, Any] | None:
"""Return a payload ``submit_answer`` can parse, or ``None`` to ignore.
Precedence (the order MUST hold):
1. **Explicit id** — if :meth:`SlackTransport.parse_answer` can already
recover a ``question_id`` (callback_id / nested view-or-message
metadata / bare ``question_id``), pass the original payload straight
through unchanged. This covers block_actions, view_submission, slash
commands, and any synthetic-but-explicit reply.
2. **Thread-reply fallback** — only when (1) raises ``ValueError`` (no
recoverable id): a real free-text Events API reply. Resolve the
question by the inner event's ``thread_ts`` against the OPEN ledger
row whose ``channel_ref`` equals it (the question message's ``ts``),
and synthesize an explicit ``{"question_id", "answer"}`` payload whose
answer is the inner ``event.text`` (stripped). The
``status='open'``-constrained lookup is anti-replay: a thread_ts for
a closed / answered / expired row resolves to ``None`` → ignored.
Returns the payload to submit, or ``None`` when no question can be
resolved (the caller ignores the event without crashing). Re-raises
nothing: a genuinely unrecoverable explicit payload surfaces as the
``ValueError`` from the final ``submit_answer`` call in ``handle_event``.
"""
try:
self._transport.parse_answer(raw_payload)
except ValueError:
pass
else:
# Explicit id recovered; submit the original payload unchanged.
return raw_payload
# No explicit id: try the events-API thread-reply mapping.
event = _inner_event(raw_payload)
thread_ts = event.get("thread_ts")
if not thread_ts:
# Not a thread reply (or no inner event): nothing to map. Returning
# the original payload lets submit_answer raise the canonical
# "no recoverable question_id" ValueError, which handle_event logs
# and ignores.
return raw_payload
question_id = find_open_question_by_channel_ref(conn, str(thread_ts))
if question_id is None:
# thread_ts matched no OPEN row (stale / replayed / answered): ignore.
_LOG.debug(
"ignoring Slack thread reply: thread_ts %r matched no open "
"question (channel_ref)",
thread_ts,
)
return None
text = event.get("text")
answer = text.strip() if isinstance(text, str) else text
# Synthesize an explicit payload the transport already understands
# (bare question_id + answer). The answer is opaque DATA — stored
# verbatim as JSON by the responder, never interpreted.
return {"question_id": question_id, "answer": answer}
def serve(self) -> None: # pragma: no cover - live socket, not unit-tested
"""Open the Socket Mode connection and forward events to ``handle_event``.
@ -372,6 +460,33 @@ class SlackListener:
_LOG.debug("SlackListener.close: handler teardown raised; ignoring")
def _inner_event(raw_payload: Mapping[str, Any]) -> Mapping[str, Any]:
"""Return the Events API inner event, or an empty mapping if there is none.
Real slack_bolt delivers an Events API message/mention as
``{"type": "event_callback", "event": {...}}``; the answer-bearing fields
(``type``, ``user``, ``text``, ``thread_ts``, ``ts``, ``channel``) all live
on the inner ``event``. Interactive / slash payloads have no inner
``event`` and yield ``{}`` here (their fields are top-level).
"""
event = raw_payload.get("event")
return event if isinstance(event, Mapping) else {}
def _discriminating_type(raw_payload: Mapping[str, Any]) -> Any:
"""Return the type used to gate an event against ``_ANSWER_BEARING_TYPES``.
For the real Events API envelope (top-level ``type == "event_callback"``)
the discriminating type is the INNER ``event.type`` (``message`` /
``app_mention``). For interactive / slash payloads it is the top-level
``type``. Returns ``None`` when no type is present (handled as a pass).
"""
top = raw_payload.get("type")
if top == _EVENT_CALLBACK_TYPE:
return _inner_event(raw_payload).get("type")
return top
def _extract_sender_id(raw_payload: Mapping[str, Any]) -> str | None:
"""Recover the inbound sender's Slack user id from any supported shape.

View file

@ -660,6 +660,14 @@ def _cmd_start(args: argparse.Namespace, *, out: Any) -> int:
either way.
"""
coordinator = _build_coordinator(args)
# start_task runs the clarifier graph to the first human gate IN THIS PROCESS,
# and the clarifier calls Claude (assess_confidence). The invoker is a
# process-local binding, so a running daemon does not help this CLI process —
# bind the real subscription invoker here, mirroring Coordinator.serve(), or
# start_task fails with "claude_invoke has no invoker bound".
from agent_team.invoker import bind_subscription_invoker
bind_subscription_invoker()
coordinator.setup()
thread_id = coordinator.start_task(
task_text=args.task, transport_name=args.transport

View file

@ -61,19 +61,31 @@ class RecordingQueue:
self.jobs.append(job)
# The default channel_ref (== the question message ``ts``) stored on a seeded
# open row. A real thread reply's ``thread_ts`` equals this value.
SEED_CHANNEL_REF = "1700000000.000100"
def _seed_open_question(
db_path: Path,
*,
thread_id: str = "t1",
question_id: str = "q1",
turn: int = 0,
channel_ref: str = SEED_CHANNEL_REF,
) -> None:
"""Insert a real ``open`` ledger row via notify_question (no network post)."""
"""Insert a real ``open`` ledger row via notify_question (no network post).
The injected poster returns ``channel_ref`` as the message ``ts``, which
``notify_question`` persists as the row's ``channel_ref``. A real Events API
thread reply carries that same value as its ``thread_ts``, so the
channel_ref → open-question mapping can resolve it.
"""
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"}
channel="C123", poster=lambda _msg: {"ts": channel_ref}
)
qs = QuestionSet(
thread_id=thread_id,
@ -87,6 +99,35 @@ def _seed_open_question(
conn.close()
def _events_api_reply(
*,
text: str = "approve",
sender_id: str | None = OWNER_ID,
thread_ts: str = SEED_CHANNEL_REF,
channel: str = "C123",
) -> dict[str, Any]:
"""A REAL slack_bolt Events API thread-reply envelope (free-text answer).
This is the shape slack_bolt actually delivers for a thread reply / message:
``{"type": "event_callback", "event": {"type": "message", ...}}`` — the
discriminating type, sender, text, and thread anchor all live on the INNER
``event``, and there is NO top-level ``callback_id`` / ``question_id`` /
metadata. The question is resolved by ``thread_ts`` (== the question
message's ``ts`` == the row's ``channel_ref``). Pass ``sender_id=None`` to
omit the author entirely (the no-recoverable-sender case).
"""
event: dict[str, Any] = {
"type": "message",
"text": text,
"thread_ts": thread_ts,
"ts": "1700000001.000200",
"channel": channel,
}
if sender_id is not None:
event["user"] = sender_id
return {"type": "event_callback", "event": event}
def _interactive_payload(
question_id: str,
value: str = "approve",
@ -279,43 +320,131 @@ def test_handle_event_rejects_payload_without_sender_id(db_path: Path) -> None:
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)
def test_handle_event_accepts_real_events_api_thread_reply(db_path: Path) -> None:
"""A REAL slack_bolt thread reply (no callback_id) maps via channel_ref.
This is the regression guard for the live-Slack inbound bug: the prior
synthetic fixture cheated with a top-level ``type:"message"`` + explicit
``callback_id`` that real Slack never sends. A real free-text reply arrives
as ``event_callback`` wrapping an inner ``message`` event with only a
``thread_ts`` to anchor it. The listener must (a) gate on the INNER
``event.type``, (b) authorize ``event.user``, and (c) resolve the question
by ``thread_ts`` == the seeded row's ``channel_ref``.
"""
_seed_open_question(
db_path, question_id="q1", thread_id="t1", turn=0, channel_ref=SEED_CHANNEL_REF
)
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)
outcome = listener.handle_event(
_events_api_reply(text="approve", thread_ts=SEED_CHANNEL_REF)
)
assert outcome is not None and outcome.accepted is True
assert 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
# The answer value is the inner event.text, stripped.
assert job.answer == "approve"
assert _row_status(db_path, "q1") == "answered"
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")
def test_handle_event_strips_thread_reply_text(db_path: Path) -> None:
"""The free-text answer is the inner ``event.text``, whitespace-stripped."""
_seed_open_question(db_path, question_id="q1", channel_ref=SEED_CHANNEL_REF)
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)
listener.handle_event(
_events_api_reply(text=" use the release branch ", thread_ts=SEED_CHANNEL_REF)
)
assert queue.jobs[0].answer == "use the release branch"
def test_handle_event_rejects_real_events_api_non_owner(db_path: Path) -> None:
"""A real Events API thread reply from a non-owner ``event.user`` is rejected.
Authorization runs FIRST on the inner ``event.user``; the channel_ref
mapping never resolves which question because submit_answer is never reached.
The seeded open row stays ``open``.
"""
_seed_open_question(db_path, question_id="q1", channel_ref=SEED_CHANNEL_REF)
queue = RecordingQueue()
listener = _listener(db_path, queue, owner_ids={OWNER_ID})
outcome = listener.handle_event(
_events_api_reply(sender_id="U_INTRUDER", thread_ts=SEED_CHANNEL_REF)
)
assert outcome is None
assert queue.jobs == []
assert _row_status(db_path, "q1") == "open"
def test_handle_event_thread_reply_no_matching_open_row_is_noop(db_path: Path) -> None:
"""A real reply whose thread_ts matches no open row is ignored (no-op)."""
# Seed an open question with a DIFFERENT channel_ref so the reply's thread_ts
# anchors to nothing open.
_seed_open_question(db_path, question_id="q1", channel_ref="9999999999.000999")
queue = RecordingQueue()
listener = _listener(db_path, queue, owner_ids={OWNER_ID})
outcome = listener.handle_event(_events_api_reply(thread_ts="1700000000.unmatched"))
assert outcome is None
assert queue.jobs == []
assert _row_status(db_path, "q1") == "open"
def test_handle_event_thread_reply_to_answered_row_is_noop(db_path: Path) -> None:
"""A reply whose thread_ts matches an ALREADY-answered row is a no-op.
The channel_ref lookup is constrained to ``status='open'``, so once the row
is answered a later thread reply on the same thread_ts resolves to nothing
and never re-opens or overwrites the answer (anti-replay).
"""
_seed_open_question(db_path, question_id="q1", channel_ref=SEED_CHANNEL_REF)
queue = RecordingQueue()
listener = _listener(db_path, queue, owner_ids={OWNER_ID})
first = listener.handle_event(
_events_api_reply(text="approve", thread_ts=SEED_CHANNEL_REF)
)
assert first is not None and first.accepted is True
queue.jobs.clear()
# Row is now 'answered'; a second reply on the same thread is a no-op.
second = listener.handle_event(
_events_api_reply(text="changed my mind", thread_ts=SEED_CHANNEL_REF)
)
assert second is None
assert queue.jobs == []
assert _row_status(db_path, "q1") == "answered"
def test_handle_event_app_mention_envelope_maps_via_thread_ts(db_path: Path) -> None:
"""An ``app_mention`` Events API reply is normalized + mapped identically.
The same inner-event normalization that fixes ``message`` also fixes
``app_mention`` (both arrive wrapped in ``event_callback``).
"""
_seed_open_question(db_path, question_id="q1", channel_ref=SEED_CHANNEL_REF)
queue = RecordingQueue()
listener = _listener(db_path, queue, owner_ids={OWNER_ID})
payload = _events_api_reply(text="<@U_BOT> approve", thread_ts=SEED_CHANNEL_REF)
payload["event"]["type"] = "app_mention"
outcome = listener.handle_event(payload)
assert outcome is not None and outcome.accepted is True
assert len(queue.jobs) == 1
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)
@ -382,14 +511,25 @@ def test_handle_event_ignores_unrelated_event_type(db_path: Path) -> None:
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)
"""An answer-bearing type with no recoverable question_id is logged-ignored.
# A message with no callback_id / metadata / question_id: parse_answer raises
# ValueError, which handle_event swallows.
The sender IS an authorized owner (so it clears AUTHZ-01) and the message is
NOT a thread reply (no ``thread_ts``), so neither explicit id recovery nor
the channel_ref fallback can resolve a question: submit_answer raises
ValueError, which handle_event swallows.
"""
queue = RecordingQueue()
listener = _listener(db_path, queue, owner_ids={OWNER_ID})
# A top-level message (no inner event, no thread_ts, no callback_id) from an
# authorized owner: nothing maps it to a question.
outcome = listener.handle_event(
{"type": "message", "text": "just chatting", "channel": "C123"}
{
"type": "message",
"text": "just chatting",
"channel": "C123",
"user": OWNER_ID,
}
)
assert outcome is None