diff --git a/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index eee73b0..2f89fe1 100644 --- a/agent-team/agent_team/coordinator.py +++ b/agent-team/agent_team/coordinator.py @@ -50,7 +50,9 @@ committed leaves (responder, resume_worker, db.schema, graph, transport). from __future__ import annotations import logging +import os import queue +import threading from datetime import timedelta from pathlib import Path 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.resume_worker import ResumeResult, ResumeWorker from agent_team.transport.base import Transport +from agent_team.transport.slack_adapter import SlackTransport if TYPE_CHECKING: # pragma: no cover - typing only from agent_team.task_model import PipelineState @@ -68,6 +71,7 @@ __all__ = [ "Coordinator", "build_verify_wiring", "default_clarify_node_factory", + "default_slack_listener_factory", "gated_build_verify_wiring", ] @@ -124,6 +128,48 @@ CheckpointerFactory = Callable[[Path], Any] # so the deadline policy stays I/O-free in tests (§6.6). 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]: """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, deadline_window: timedelta | None = None, alarm_hook: AlarmHook | None = None, + build_listener: ListenerFactory | None = None, ) -> None: self._db_path = Path(db_path) self._transport = transport @@ -357,6 +404,11 @@ class Coordinator: self._resume_queue: "queue.Queue[Any]" = resume_queue or queue.Queue() self._deadline_window = deadline_window or graph_mod.DEFAULT_CLARIFY_DEADLINE 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(). self._graph: Any = None @@ -364,6 +416,10 @@ class Coordinator: # Retains a context-manager checkpointer (production SQLite saver) so its # __exit__ is not run early; held open for the daemon's lifetime. 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). @@ -814,32 +870,147 @@ class Coordinator: # ------------------------------------------------------------------ # 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 (:func:`agent_team.invoker.bind_subscription_invoker`) BEFORE :meth:`setup` builds the clarifier node (so the node's ``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. - The actual Slack inbound feed is the slack_listener's job; the - coordinator exposes :meth:`submit_answer` and the shared - :attr:`resume_queue` for it. This loop owns only the deadline/recovery - maintenance cadence. + The Slack inbound feed (Socket Mode) is the slack_listener's job and runs + concurrently on a background daemon thread; this loop owns the + deadline/recovery maintenance cadence. On shutdown (Ctrl-C / + 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 bind_subscription_invoker() self.setup() self.recover() + self._maybe_start_slack_listener() interval = (poll_interval or DEFAULT_POLL_INTERVAL).total_seconds() import time - while True: # pragma: no cover - the infinite daemon loop - self.tick() - time.sleep(interval) + try: + while True: # pragma: no cover - the infinite daemon loop + 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. diff --git a/agent-team/agent_team/transport/slack_listener.py b/agent-team/agent_team/transport/slack_listener.py index b592a56..d0472dd 100644 --- a/agent-team/agent_team/transport/slack_listener.py +++ b/agent-team/agent_team/transport/slack_listener.py @@ -147,6 +147,10 @@ class SlackListener: # 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() + # 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: """Normalize + submit one inbound event; return its outcome or ``None``. @@ -331,7 +335,35 @@ class SlackListener: def _on_mention(body: Mapping[str, Any]) -> None: _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: diff --git a/agent-team/tests/test_coordinator.py b/agent-team/tests/test_coordinator.py index d0d67c5..47b2601 100644 --- a/agent-team/tests/test_coordinator.py +++ b/agent-team/tests/test_coordinator.py @@ -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 finally: 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 diff --git a/agent-team/tests/test_slack_listener.py b/agent-team/tests/test_slack_listener.py index 10eaf85..f1d04ff 100644 --- a/agent-team/tests/test_slack_listener.py +++ b/agent-team/tests/test_slack_listener.py @@ -449,6 +449,36 @@ def test_serve_requires_tokens(db_path: Path) -> None: 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: """Sanity: the injected SlackTransport is the contract the listener expects.""" assert isinstance(SlackTransport(channel="C123"), Transport)