fix(agent-team): serve() starts the inbound Slack listener (D-1)

Coordinator.serve() now constructs and starts the SlackListener concurrently
with the tick/drain loop on a background daemon thread, but ONLY when the live
transport is a SlackTransport AND SLACK_APP_TOKEN is configured. When Slack is
not the transport or the app token is absent, serve() behaves exactly as before
(tick/recover only) — Slack is never made mandatory.

- New injectable build_listener seam + default_slack_listener_factory sharing
  the coordinator's own transport, ledger db_path, and resume_queue put.
- AUTHZ-01 owner-allowlist + open-status CAS untouched: serve() sources
  AGENT_TEAM_SLACK_OWNER_IDS in SlackListener.serve, which still fails closed.
- SlackListener.close() added for clean Socket Mode teardown on shutdown;
  serve() stops the listener + joins the thread in a finally.
- Tests: start-when-Slack+app-token, no-start otherwise, clean shutdown,
  idempotent start, serve start/stop around the loop, listener close().
This commit is contained in:
Adam Moussa 2026-06-18 16:39:32 -04:00
parent 95e481890a
commit 06811113b9
4 changed files with 411 additions and 10 deletions

View file

@ -50,7 +50,9 @@ committed leaves (responder, resume_worker, db.schema, graph, transport).
from __future__ import annotations from __future__ import annotations
import logging import logging
import os
import queue import queue
import threading
from datetime import timedelta from datetime import timedelta
from pathlib import Path from pathlib import Path
from typing import TYPE_CHECKING, Any, Callable from typing import TYPE_CHECKING, Any, Callable
@ -60,6 +62,7 @@ from agent_team import responder as responder_mod
from agent_team.db.schema import connect, init_db from agent_team.db.schema import connect, init_db
from agent_team.resume_worker import ResumeResult, ResumeWorker from agent_team.resume_worker import ResumeResult, ResumeWorker
from agent_team.transport.base import Transport from agent_team.transport.base import Transport
from agent_team.transport.slack_adapter import SlackTransport
if TYPE_CHECKING: # pragma: no cover - typing only if TYPE_CHECKING: # pragma: no cover - typing only
from agent_team.task_model import PipelineState from agent_team.task_model import PipelineState
@ -68,6 +71,7 @@ __all__ = [
"Coordinator", "Coordinator",
"build_verify_wiring", "build_verify_wiring",
"default_clarify_node_factory", "default_clarify_node_factory",
"default_slack_listener_factory",
"gated_build_verify_wiring", "gated_build_verify_wiring",
] ]
@ -124,6 +128,48 @@ CheckpointerFactory = Callable[[Path], Any]
# so the deadline policy stays I/O-free in tests (§6.6). # so the deadline policy stays I/O-free in tests (§6.6).
AlarmHook = Callable[[str], None] AlarmHook = Callable[[str], None]
# A listener factory: builds the inbound Slack Socket Mode listener
# (:class:`~agent_team.transport.slack_listener.SlackListener`) the serve loop
# starts on a background thread when Slack is the live transport and the app
# token is present. Injected so :meth:`Coordinator.serve` is testable with a fake
# listener and no live socket. The default
# (:func:`default_slack_listener_factory`) builds the real listener from the
# coordinator's shared transport / db / resume-queue and the Slack env tokens.
ListenerFactory = Callable[[], Any]
def default_slack_listener_factory(
*,
transport: SlackTransport,
db_path: Path,
enqueue_resume: Callable[[Any], None],
) -> Any:
"""Build the live :class:`SlackListener` from the coordinator's seams (D-1).
The production listener shares the coordinator's *own* live ``SlackTransport``
(so ``parse_answer`` matches the outbound ``post_question`` wiring), the same
durable ledger ``db_path``, and the same ``resume_queue`` ``put`` callable
(the documented slack_listener ⇄ drainer handoff seam — ``coordinator.py``
module docstring §3.3). The app/bot tokens and the owner allowlist are sourced
from the environment by :meth:`SlackListener.serve` itself
(``SLACK_APP_TOKEN`` / ``SLACK_BOT_TOKEN`` / ``AGENT_TEAM_SLACK_OWNER_IDS``);
the allowlist still FAILS CLOSED if ``AGENT_TEAM_SLACK_OWNER_IDS`` is unset
(AUTHZ-01), so this factory deliberately does not weaken that — it injects no
``owner_ids`` and lets ``serve`` read + enforce them.
Imported lazily for the same import-hygiene reason as the clarifier / planner
factories (the listener pulls the transport + responder leaves).
"""
from agent_team.transport.slack_listener import SlackListener
return SlackListener(
transport,
db_path,
enqueue_resume,
app_token=os.environ.get("SLACK_APP_TOKEN") or None,
bot_token=os.environ.get("SLACK_BOT_TOKEN") or None,
)
def default_clarify_node_factory() -> Callable[[PipelineState], PipelineState]: def default_clarify_node_factory() -> Callable[[PipelineState], PipelineState]:
"""Build the live Claude-backed clarifier node (§3.3, §7.1 P1). """Build the live Claude-backed clarifier node (§3.3, §7.1 P1).
@ -338,6 +384,7 @@ class Coordinator:
resume_queue: "queue.Queue[Any] | None" = None, resume_queue: "queue.Queue[Any] | None" = None,
deadline_window: timedelta | None = None, deadline_window: timedelta | None = None,
alarm_hook: AlarmHook | None = None, alarm_hook: AlarmHook | None = None,
build_listener: ListenerFactory | None = None,
) -> None: ) -> None:
self._db_path = Path(db_path) self._db_path = Path(db_path)
self._transport = transport self._transport = transport
@ -357,6 +404,11 @@ class Coordinator:
self._resume_queue: "queue.Queue[Any]" = resume_queue or queue.Queue() self._resume_queue: "queue.Queue[Any]" = resume_queue or queue.Queue()
self._deadline_window = deadline_window or graph_mod.DEFAULT_CLARIFY_DEADLINE self._deadline_window = deadline_window or graph_mod.DEFAULT_CLARIFY_DEADLINE
self._alarm_hook = alarm_hook or self._default_alarm_hook self._alarm_hook = alarm_hook or self._default_alarm_hook
# The inbound Slack listener factory (D-1). Left None, serve() uses the
# default that builds the real SlackListener; tests inject a fake to prove
# the start/no-start/shutdown wiring with no live socket. Only consulted
# by serve() when Slack is the live transport AND the app token is set.
self._build_listener = build_listener
# Built by setup(). # Built by setup().
self._graph: Any = None self._graph: Any = None
@ -364,6 +416,10 @@ class Coordinator:
# Retains a context-manager checkpointer (production SQLite saver) so its # Retains a context-manager checkpointer (production SQLite saver) so its
# __exit__ is not run early; held open for the daemon's lifetime. # __exit__ is not run early; held open for the daemon's lifetime.
self._checkpointer_cm: Any = None self._checkpointer_cm: Any = None
# The running inbound listener + its daemon thread (None until serve()
# starts one). Held so serve()'s shutdown can stop it cleanly.
self._listener: Any = None
self._listener_thread: threading.Thread | None = None
# ------------------------------------------------------------------ # # ------------------------------------------------------------------ #
# Accessors (the shared queue is the slack_listener handoff seam). # Accessors (the shared queue is the slack_listener handoff seam).
@ -814,32 +870,147 @@ class Coordinator:
# ------------------------------------------------------------------ # # ------------------------------------------------------------------ #
def serve(self, *, poll_interval: timedelta | None = None) -> None: def serve(self, *, poll_interval: timedelta | None = None) -> None:
"""Run the live daemon: bind invoker, setup, recover, then tick forever. """Run the live daemon: bind invoker, setup, recover, start listener, tick.
The production entry. Binds the real Claude invoker The production entry. Binds the real Claude invoker
(:func:`agent_team.invoker.bind_subscription_invoker`) BEFORE (:func:`agent_team.invoker.bind_subscription_invoker`) BEFORE
:meth:`setup` builds the clarifier node (so the node's :meth:`setup` builds the clarifier node (so the node's
``billing.claude_invoke`` calls hit the live subscription path), runs the ``billing.claude_invoke`` calls hit the live subscription path), runs the
startup :meth:`recover` sweep, then loops calling :meth:`tick` on the startup :meth:`recover` sweep, **starts the inbound Slack listener**
(:meth:`_maybe_start_slack_listener`) when Slack is the live transport
and the app token is configured, then loops calling :meth:`tick` on the
deadline cadence. deadline cadence.
The actual Slack inbound feed is the slack_listener's job; the The Slack inbound feed (Socket Mode) is the slack_listener's job and runs
coordinator exposes :meth:`submit_answer` and the shared concurrently on a background daemon thread; this loop owns the
:attr:`resume_queue` for it. This loop owns only the deadline/recovery deadline/recovery maintenance cadence. On shutdown (Ctrl-C /
maintenance cadence. SIGTERM-driven ``KeyboardInterrupt``) the listener is stopped cleanly via
:meth:`_stop_slack_listener` in a ``finally``.
Slack is NOT mandatory: when the transport is not the live ``SlackTransport``
or the app token is absent, no listener is started and ``serve`` behaves
exactly as before (tick/recover only).
""" """
from agent_team.invoker import bind_subscription_invoker from agent_team.invoker import bind_subscription_invoker
bind_subscription_invoker() bind_subscription_invoker()
self.setup() self.setup()
self.recover() self.recover()
self._maybe_start_slack_listener()
interval = (poll_interval or DEFAULT_POLL_INTERVAL).total_seconds() interval = (poll_interval or DEFAULT_POLL_INTERVAL).total_seconds()
import time import time
while True: # pragma: no cover - the infinite daemon loop try:
self.tick() while True: # pragma: no cover - the infinite daemon loop
time.sleep(interval) self.tick()
time.sleep(interval)
finally:
self._stop_slack_listener()
# ------------------------------------------------------------------ #
# Inbound Slack listener wiring (D-1): start concurrently with the
# tick/drain loop ONLY when Slack is the live transport and the app
# token is configured; stop cleanly on shutdown.
# ------------------------------------------------------------------ #
def _slack_listener_enabled(self) -> bool:
"""Return ``True`` iff the inbound Slack listener should run (D-1 gate).
Two conditions, both required (and Slack must stay OPTIONAL):
* the live transport is a :class:`SlackTransport` (the inbound
``parse_answer`` must match the outbound poster); and
* ``SLACK_APP_TOKEN`` is present in the environment (Socket Mode needs the
app-level token to open the outbound WebSocket).
If either is false the daemon runs WITHOUT an inbound listener exactly as
before — Slack is never made mandatory. The owner allowlist
(``AGENT_TEAM_SLACK_OWNER_IDS``) is intentionally NOT part of this gate:
the listener fails closed on an empty allowlist (AUTHZ-01), so starting it
unprovisioned safely rejects every answer rather than silently never
hearing Slack.
"""
if not isinstance(self._transport, SlackTransport):
return False
return bool(os.environ.get("SLACK_APP_TOKEN"))
def _maybe_start_slack_listener(self) -> bool:
"""Start the inbound Slack listener on a daemon thread if enabled (D-1).
Builds the listener via the injected ``build_listener`` factory (or the
:func:`default_slack_listener_factory`, sharing the coordinator's own live
transport, durable ``db_path`` and ``resume_queue`` put), then runs its
blocking :meth:`~agent_team.transport.slack_listener.SlackListener.serve`
on a background **daemon** thread so the Socket Mode socket and the
tick/drain loop run concurrently. Returns ``True`` iff a listener was
started; ``False`` (and no thread) when :meth:`_slack_listener_enabled`
is false — the Slack-optional contract.
Idempotent: a second call while a listener thread is already alive is a
no-op.
"""
if self._listener_thread is not None and self._listener_thread.is_alive():
return True
if not self._slack_listener_enabled():
_LOG.info(
"inbound Slack listener not started (transport is not live Slack "
"or SLACK_APP_TOKEN is unset); daemon runs tick/recover only"
)
return False
if self._build_listener is not None:
listener = self._build_listener()
else:
listener = default_slack_listener_factory(
transport=self._transport,
db_path=self._db_path,
enqueue_resume=self._resume_queue.put,
)
self._listener = listener
thread = threading.Thread(
target=self._run_listener,
name="agent-team-slack-listener",
daemon=True,
)
self._listener_thread = thread
thread.start()
_LOG.info("inbound Slack listener started (Socket Mode, background thread)")
return True
def _run_listener(self) -> None:
"""Thread target: run the listener's blocking serve, log on exit.
Swallows any listener exception into a log line so a listener crash takes
down only the inbound socket (which Restart=on-failure / a manual restart
recovers), never the maintenance loop or the whole process from a
background thread.
"""
try:
self._listener.serve()
except Exception: # noqa: BLE001 - isolate the listener thread
_LOG.exception("inbound Slack listener thread exited with an error")
def _stop_slack_listener(self) -> None:
"""Stop the inbound listener cleanly on daemon shutdown (D-1).
Best-effort: closes the Socket Mode connection via the listener's
``close`` (if it exposes one), then joins the daemon thread briefly. A
teardown error never propagates — shutdown must complete regardless.
"""
listener = self._listener
thread = self._listener_thread
self._listener = None
self._listener_thread = None
if listener is not None:
close = getattr(listener, "close", None)
if close is not None:
try:
close()
except Exception: # noqa: BLE001 - shutdown must not raise
_LOG.debug("listener close raised during shutdown; ignoring")
if thread is not None and thread.is_alive():
thread.join(timeout=5.0)
# ------------------------------------------------------------------ # # ------------------------------------------------------------------ #
# Internals. # Internals.

View file

@ -147,6 +147,10 @@ class SlackListener:
# The owner allowlist (AUTHZ-01). An empty set is the fail-closed default: # The owner allowlist (AUTHZ-01). An empty set is the fail-closed default:
# an unconfigured deploy rejects every answer. # an unconfigured deploy rejects every answer.
self._owner_ids: set[str] = set(owner_ids) if owner_ids else set() self._owner_ids: set[str] = set(owner_ids) if owner_ids else set()
# The live Socket Mode handler, retained by :meth:`serve` so :meth:`close`
# can stop it cleanly on daemon shutdown. ``None`` until ``serve`` opens
# the socket.
self._socket_handler: Any = None
def handle_event(self, raw_payload: Any) -> AnswerOutcome | None: def handle_event(self, raw_payload: Any) -> AnswerOutcome | None:
"""Normalize + submit one inbound event; return its outcome or ``None``. """Normalize + submit one inbound event; return its outcome or ``None``.
@ -331,7 +335,35 @@ class SlackListener:
def _on_mention(body: Mapping[str, Any]) -> None: def _on_mention(body: Mapping[str, Any]) -> None:
_forward(body) _forward(body)
SocketModeHandler(app, self._app_token).start() handler = SocketModeHandler(app, self._app_token)
self._socket_handler = handler
handler.start()
def close(self) -> None:
"""Best-effort clean stop of the Socket Mode connection (daemon shutdown).
The coordinator's :meth:`~agent_team.coordinator.Coordinator.serve` calls
this when it shuts the daemon down so the outbound WebSocket is closed
cleanly rather than only dying with the process. Tolerant: if the handler
was never opened, or the SDK's ``close``/``disconnect`` raises, the error
is swallowed — shutdown must never be blocked by a transport teardown
failure. The handler is dropped afterward so a second ``close`` is a
no-op.
"""
handler = self._socket_handler
if handler is None:
return
self._socket_handler = None
# slack_bolt's SocketModeHandler exposes ``close`` (and the underlying
# client a ``disconnect``); try the most specific available, swallow any
# teardown error.
closer = getattr(handler, "close", None) or getattr(handler, "disconnect", None)
if closer is None:
return
try:
closer()
except Exception: # noqa: BLE001 - shutdown must not raise
_LOG.debug("SlackListener.close: handler teardown raised; ignoring")
def _extract_sender_id(raw_payload: Mapping[str, Any]) -> str | None: def _extract_sender_id(raw_payload: Mapping[str, Any]) -> str | None:

View file

@ -602,3 +602,171 @@ def test_setup_with_p2_factories_builds_a_review_node(db_path: Path) -> None:
assert graph_mod.REVIEW in coord.graph.get_graph().nodes assert graph_mod.REVIEW in coord.graph.get_graph().nodes
finally: finally:
review_loop._review_invoker = saved review_loop._review_invoker = saved
# --------------------------------------------------------------------------- #
# Inbound Slack listener wiring in serve() (D-1)
# --------------------------------------------------------------------------- #
class _FakeListener:
"""A record-only stand-in for SlackListener (no socket, no SDK).
``serve`` blocks on an Event until ``close`` is called, mirroring the live
listener whose ``serve`` blocks on the Socket Mode handler until shut down.
Records that ``serve`` ran and that ``close`` was called so the coordinator
start/shutdown wiring can be asserted.
"""
def __init__(self) -> None:
import threading as _t
self.served = False
self.closed = False
self._stop = _t.Event()
def serve(self) -> None:
self.served = True
self._stop.wait(timeout=5.0)
def close(self) -> None:
self.closed = True
self._stop.set()
def _slack_coordinator(
db_path: Path,
*,
build_listener: Any = None,
) -> Coordinator:
"""A Coordinator whose live transport IS a SlackTransport (listener gate)."""
from agent_team.transport.slack_adapter import SlackTransport
saver = _Saver()
return Coordinator(
db_path=db_path,
transport=SlackTransport("C123"),
build_clarify_node=lambda: graph_mod.clarify_node,
build_checkpointer=lambda _path: saver,
build_listener=build_listener,
)
def test_serve_starts_listener_when_slack_and_app_token(
db_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
"""serve() starts the inbound listener when Slack is the transport AND the
app token is configured, then stops it cleanly on shutdown."""
monkeypatch.setenv("SLACK_APP_TOKEN", "xapp-test")
listener = _FakeListener()
coord = _slack_coordinator(db_path, build_listener=lambda: listener)
coord.setup()
assert coord._maybe_start_slack_listener() is True
# The background daemon thread actually ran the listener's serve().
import time
for _ in range(50):
if listener.served:
break
time.sleep(0.01)
assert listener.served is True
# Clean shutdown stops the listener and joins the thread.
coord._stop_slack_listener()
assert listener.closed is True
assert coord._listener is None
assert coord._listener_thread is None
def test_serve_does_not_start_listener_without_app_token(
db_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
"""No app token => no inbound listener (Slack stays optional)."""
monkeypatch.delenv("SLACK_APP_TOKEN", raising=False)
listener = _FakeListener()
coord = _slack_coordinator(db_path, build_listener=lambda: listener)
coord.setup()
assert coord._slack_listener_enabled() is False
assert coord._maybe_start_slack_listener() is False
assert listener.served is False
assert coord._listener is None
assert coord._listener_thread is None
def test_serve_does_not_start_listener_when_transport_not_slack(
db_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
"""A non-Slack transport never starts the listener even with the app token."""
monkeypatch.setenv("SLACK_APP_TOKEN", "xapp-test")
listener = _FakeListener()
saver = _Saver()
coord = Coordinator(
db_path=db_path,
transport=FakeTransport(), # not a SlackTransport
build_clarify_node=lambda: graph_mod.clarify_node,
build_checkpointer=lambda _path: saver,
build_listener=lambda: listener,
)
coord.setup()
assert coord._slack_listener_enabled() is False
assert coord._maybe_start_slack_listener() is False
assert listener.served is False
def test_maybe_start_listener_is_idempotent(
db_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
"""A second start while the thread is alive is a no-op (one listener)."""
monkeypatch.setenv("SLACK_APP_TOKEN", "xapp-test")
built: list[_FakeListener] = []
def _factory() -> _FakeListener:
lst = _FakeListener()
built.append(lst)
return lst
coord = _slack_coordinator(db_path, build_listener=_factory)
coord.setup()
try:
assert coord._maybe_start_slack_listener() is True
assert coord._maybe_start_slack_listener() is True
assert len(built) == 1 # not rebuilt/restarted
finally:
coord._stop_slack_listener()
def test_stop_listener_is_safe_when_none(db_path: Path) -> None:
"""Stopping with no listener running never raises (clean-shutdown safety)."""
coord = _slack_coordinator(db_path)
coord._stop_slack_listener() # must be a quiet no-op
assert coord._listener is None
def test_serve_wires_start_and_stop_around_the_loop(
db_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
"""serve() starts the listener before the loop and stops it in finally even
when the loop is interrupted (KeyboardInterrupt)."""
monkeypatch.setenv("SLACK_APP_TOKEN", "xapp-test")
listener = _FakeListener()
coord = _slack_coordinator(db_path, build_listener=lambda: listener)
# Neuter the live invoker bind so serve() needs no Claude SDK / token.
# serve() does `from agent_team.invoker import bind_subscription_invoker`,
# so patch the name on the invoker module it imports from.
import agent_team.invoker as invoker_mod
monkeypatch.setattr(invoker_mod, "bind_subscription_invoker", lambda: None)
# Break the infinite loop on the first tick.
monkeypatch.setattr(
coord, "tick", lambda: (_ for _ in ()).throw(KeyboardInterrupt())
)
with pytest.raises(KeyboardInterrupt):
coord.serve(poll_interval=timedelta(seconds=0))
assert listener.served is True # started before the loop
assert listener.closed is True # stopped in finally on interrupt

View file

@ -449,6 +449,36 @@ def test_serve_requires_tokens(db_path: Path) -> None:
listener.serve() listener.serve()
def test_close_without_open_handler_is_noop(db_path: Path) -> None:
"""close() before serve() ever opened a socket is a quiet no-op."""
listener = _listener(db_path, RecordingQueue())
listener.close() # must not raise
listener.close() # idempotent
def test_close_stops_socket_handler(db_path: Path) -> None:
"""close() invokes the retained Socket Mode handler's close and drops it."""
class _FakeHandler:
def __init__(self) -> None:
self.closed = False
def close(self) -> None:
self.closed = True
listener = _listener(db_path, RecordingQueue())
handler = _FakeHandler()
listener._socket_handler = handler
listener.close()
assert handler.closed is True
assert listener._socket_handler is None
# A teardown that raises is swallowed (shutdown must never raise).
listener._socket_handler = type(
"_Boom", (), {"close": lambda self: (_ for _ in ()).throw(RuntimeError("x"))}
)()
listener.close() # no exception
def test_real_slack_transport_is_a_transport() -> None: def test_real_slack_transport_is_a_transport() -> None:
"""Sanity: the injected SlackTransport is the contract the listener expects.""" """Sanity: the injected SlackTransport is the contract the listener expects."""
assert isinstance(SlackTransport(channel="C123"), Transport) assert isinstance(SlackTransport(channel="C123"), Transport)