From 9caf483117b5e4549c5b1fc7583b598a464ca808 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Mon, 22 Jun 2026 15:58:01 -0400 Subject: [PATCH] fix(agent-team): harden listener respawn/close, broaden handle_event guard, channel_ref partial-unique MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Three robustness/hardening fixes surfaced by /sh-security-review on the agent-team listener/coordinator surface. None alter AUTHZ-01 allowlist behavior or the first-answer-wins compare-and-set semantics. 1. Respawn close() leak (CWE-772). The watchdog _supervise_slack_listener respawned the inbound Slack listener without tearing down the dead one, leaking a Socket Mode WebSocket / SDK thread set per flap. Now the dead listener is closed before respawn (new _close_dead_listener, idempotent), AND _run_listener has a finally that always closes the listener so a crashed serve() releases its socket. SlackListener.close() is idempotent, so the belt-and-braces close stays a safe no-op. 2. Broadened exception guard in handle_event (CWE-248). The submit block only caught ValueError; the accept path (submit_answer -> _question_turn/_question_thread) can raise KeyError on a concurrently mutated row, and the CAS can raise sqlite3.Error. An uncaught exception would escape into the Bolt dispatch. Added a separate `except Exception` that logs at WARNING (not silent, not debug) and returns None. The existing ValueError-as-debug behavior is unchanged; authorization still runs first, so the trust boundary is not widened. 3. channel_ref partial-unique index (defense-in-depth). Added uq_pending_questions_open_channel_ref — a PARTIAL UNIQUE index on (channel_ref) WHERE channel_ref IS NOT NULL AND status='open' — so two OPEN rows can never share a non-null channel_ref (a thread_ts can never map to two open questions). Installed in init_db AND unconditionally in migrate (idempotent IF NOT EXISTS) so existing v1 DBs gain it. NULLs and closed rows are excluded; mirrored verbatim into schema.sql. Tests: +8 (was 960, now 968). New: schema partial-unique reject/null/closed/ migrate cases; handle_event KeyError + sqlite3.Error swallow cases; coordinator close-before-respawn + run_listener-closes-on-crash. Fixed the operator-cli test fixture to use a per-question channel_ref (it previously inserted multiple open rows sharing one ref, which the new index correctly rejects). --- agent-team/agent_team/coordinator.py | 51 +++++++- agent-team/agent_team/db/schema.py | 16 +++ agent-team/agent_team/db/schema.sql | 7 ++ .../agent_team/transport/slack_listener.py | 17 +++ agent-team/tests/test_coordinator.py | 78 ++++++++++++ agent-team/tests/test_operator_cli.py | 8 +- agent-team/tests/test_schema.py | 113 ++++++++++++++++++ agent-team/tests/test_slack_listener.py | 65 ++++++++++ 8 files changed, 352 insertions(+), 3 deletions(-) diff --git a/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index d4f31e0..36fc6a9 100644 --- a/agent-team/agent_team/coordinator.py +++ b/agent-team/agent_team/coordinator.py @@ -1000,10 +1000,38 @@ class Coordinator: "ALARM: inbound Slack listener is DOWN — answers are NOT being " "received; clarifier gates will park. Respawning the listener." ) - # Drop the dead handle so the idempotent starter actually respawns. + # Tear down the DEAD listener before respawning (CWE-772 resource leak). + # The crashed thread's serve() may have left a Socket Mode WebSocket / + # SDK threads dangling; over repeated flaps a respawn-without-close leaks + # one socket+thread set per flap. SlackListener.close() is idempotent (it + # nulls _socket_handler first), so closing a possibly-already-torn-down + # listener here is safe. Drop both handles so the idempotent starter + # actually constructs a fresh listener. + self._close_dead_listener() + self._listener = None self._listener_thread = None self._maybe_start_slack_listener() + def _close_dead_listener(self) -> None: + """Best-effort close of a crashed listener before respawn (CWE-772). + + Releases the dead listener's Socket Mode socket / SDK threads so a + flapping listener does not leak one resource set per respawn. Tolerant: + a missing listener, a listener with no ``close``, or a ``close`` that + raises are all swallowed — supervision must never be blocked by a + teardown failure. + """ + listener = self._listener + if listener is None: + return + close = getattr(listener, "close", None) + if close is None: + return + try: + close() + except Exception: # noqa: BLE001 - respawn must not be blocked by teardown + _LOG.debug("dead listener close raised during respawn; ignoring") + def _run_listener(self) -> None: """Thread target: run the listener's blocking serve, log on exit. @@ -1011,11 +1039,30 @@ class Coordinator: 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. + + The ``finally`` ALWAYS closes the listener on the way out (CWE-772): a + crashed ``serve()`` would otherwise leave its Socket Mode WebSocket / SDK + threads dangling, leaking one resource set per flap as the supervisor + respawns. ``SlackListener.close()`` is idempotent (it nulls + ``_socket_handler`` first), so the supervisor's belt-and-braces close on + respawn stays a safe no-op after this one runs. """ + listener = self._listener try: - self._listener.serve() + listener.serve() except Exception: # noqa: BLE001 - isolate the listener thread _LOG.exception("inbound Slack listener thread exited with an error") + finally: + # Release this listener's socket/threads even on a crash. Bind the + # local above so a concurrent respawn swapping self._listener cannot + # make us close the wrong (new) listener. + if listener is not None: + close = getattr(listener, "close", None) + if close is not None: + try: + close() + except Exception: # noqa: BLE001 - teardown must not raise + _LOG.debug("listener close raised on thread exit; ignoring") def _stop_slack_listener(self) -> None: """Stop the inbound listener cleanly on daemon shutdown (D-1). diff --git a/agent-team/agent_team/db/schema.py b/agent-team/agent_team/db/schema.py index 7f60289..d2d0067 100644 --- a/agent-team/agent_team/db/schema.py +++ b/agent-team/agent_team/db/schema.py @@ -78,6 +78,9 @@ CREATE INDEX IF NOT EXISTS idx_pending_questions_thread ON pending_questions (thread_id, turn); CREATE INDEX IF NOT EXISTS idx_pending_questions_status ON pending_questions (status); +CREATE UNIQUE INDEX IF NOT EXISTS uq_pending_questions_open_channel_ref + ON pending_questions (channel_ref) + WHERE channel_ref IS NOT NULL AND status = 'open'; """.strip() BUDGET_LEDGER_DDL: str = """ @@ -212,6 +215,19 @@ def migrate(conn: sqlite3.Connection) -> None: # Future steps go here: `if current < 2: ...; current = 2`. + # Defense-in-depth index, applied UNCONDITIONALLY (idempotent IF NOT EXISTS) + # so an already-stamped v1 DB — which skips the `current < 1` block above — + # still gains the partial unique index on (channel_ref) WHERE open. This is a + # pure add-on guard (no version bump): two OPEN rows can never share a + # non-null channel_ref. Safe on existing data — NULL channel_refs are fine + # (the WHERE clause excludes them and SQLite treats NULLs as distinct), and + # there should be no existing duplicate non-null OPEN channel_refs (the + # responder records one channel_ref per posted open question, and recovery + # clears unposted ones). If a legacy DB *did* hold a duplicate, this CREATE + # would raise IntegrityError loudly rather than silently — flag, don't force. + for stmt in _split_statements(PENDING_QUESTIONS_INDEXES_DDL): + conn.execute(stmt) + conn.execute( "INSERT INTO schema_meta (id, schema_version) VALUES (1, ?) " "ON CONFLICT(id) DO UPDATE SET schema_version = excluded.schema_version", diff --git a/agent-team/agent_team/db/schema.sql b/agent-team/agent_team/db/schema.sql index 5dc4479..db4f056 100644 --- a/agent-team/agent_team/db/schema.sql +++ b/agent-team/agent_team/db/schema.sql @@ -34,6 +34,13 @@ CREATE INDEX IF NOT EXISTS idx_pending_questions_thread CREATE INDEX IF NOT EXISTS idx_pending_questions_status ON pending_questions (status); +-- Defense-in-depth: at most one OPEN question may carry a given non-null +-- channel_ref, so a thread reply's thread_ts can never map to two open rows +-- (a partial unique index; NULL channel_refs and closed rows are excluded). +CREATE UNIQUE INDEX IF NOT EXISTS uq_pending_questions_open_channel_ref + ON pending_questions (channel_ref) + WHERE channel_ref IS NOT NULL AND status = 'open'; + -- budget_ledger: the persistent shared Claude budget ledger (§6.1, §6.6). -- One row per accounted spend event; the shared daily cap and the -- interactive-first reserve are computed by summing over a UTC day. Spend is diff --git a/agent-team/agent_team/transport/slack_listener.py b/agent-team/agent_team/transport/slack_listener.py index 5a463dd..cc73b34 100644 --- a/agent-team/agent_team/transport/slack_listener.py +++ b/agent-team/agent_team/transport/slack_listener.py @@ -241,6 +241,23 @@ class SlackListener: # This is expected for unrelated chatter on the channel; ignore it. _LOG.debug("ignoring Slack event with no recoverable answer: %s", exc) return None + except Exception: # noqa: BLE001 - keep the listen loop alive on a race + # Defense-in-depth (CWE-248): the accept path runs + # submit_answer -> _question_turn / _question_thread, which raise + # KeyError if the just-answered row is concurrently mutated, and the + # underlying compare-and-set can raise sqlite3.Error under a write + # fault. An exception escaping handle_event would propagate into the + # Bolt dispatch and could take the listener down. Log-and-ignore at + # WARNING (NOT debug — this is unexpected, unlike the ValueError + # chatter case) so a real bug is still surfaced loudly, but a racing + # / malformed event cannot crash the loop. Authorization already ran + # above, so this never widens the trust boundary. + _LOG.warning( + "ignoring Slack event: unexpected error during answer submit " + "(racing row mutation / DB fault); event dropped", + exc_info=True, + ) + return None finally: conn.close() diff --git a/agent-team/tests/test_coordinator.py b/agent-team/tests/test_coordinator.py index 0b9439c..0692e35 100644 --- a/agent-team/tests/test_coordinator.py +++ b/agent-team/tests/test_coordinator.py @@ -774,6 +774,84 @@ def test_supervise_respawns_dead_listener_and_alarms( coord._stop_slack_listener() +def test_supervise_closes_dead_listener_before_respawn( + db_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + """CWE-772: the dead listener is closed BEFORE a fresh one is constructed. + + Without this, a flapping listener leaks one Socket Mode WebSocket / SDK + thread set per respawn. Record the close/build event order and assert the + dead listener's close() runs strictly before the replacement is built. + """ + monkeypatch.setenv("SLACK_APP_TOKEN", "xapp-test") + events: list[str] = [] + + class _SpyListener(_FakeListener): + def __init__(self) -> None: + super().__init__() + events.append("build") + + def close(self) -> None: + events.append("close") + super().close() + + coord = _slack_coordinator(db_path, build_listener=_SpyListener) + coord.setup() + try: + assert coord._maybe_start_slack_listener() is True + first = coord._listener + first_thread = coord._listener_thread + assert first is not None and first_thread is not None + # Kill the thread so the supervisor sees it dead, but DO NOT close the + # listener ourselves — the supervisor must be the one to close it. + first._stop.set() + first_thread.join(timeout=5.0) + assert not first_thread.is_alive() + + events.clear() + coord._supervise_slack_listener() + + # The dead listener was closed BEFORE the replacement was constructed. + assert "close" in events and "build" in events + assert events.index("close") < events.index("build") + # And the dead listener really did get closed. + assert first.closed is True + # A fresh, distinct listener is now running. + assert coord._listener is not None and coord._listener is not first + finally: + coord._stop_slack_listener() + + +def test_run_listener_closes_listener_on_crash( + db_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + """CWE-772: a crashed serve() always releases its socket via the finally. + + _run_listener's finally must close the listener even when serve() raises, so + a crash never dangles the Socket Mode socket / SDK threads. + """ + monkeypatch.setenv("SLACK_APP_TOKEN", "xapp-test") + + class _CrashingListener: + def __init__(self) -> None: + self.closed = False + + def serve(self) -> None: + raise RuntimeError("socket boom") + + def close(self) -> None: + self.closed = True + + listener = _CrashingListener() + coord = _slack_coordinator(db_path, build_listener=lambda: listener) + coord.setup() + # Drive the thread target directly (no real thread needed): the crash must be + # swallowed and close() must still run from the finally. + coord._listener = listener + coord._run_listener() + assert listener.closed is True + + def test_supervise_is_noop_when_listener_alive( db_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: diff --git a/agent-team/tests/test_operator_cli.py b/agent-team/tests/test_operator_cli.py index 4d21d2d..633198c 100644 --- a/agent-team/tests/test_operator_cli.py +++ b/agent-team/tests/test_operator_cli.py @@ -40,9 +40,15 @@ def _insert_open_question( turn: int = 0, status: str = "open", transport: str = "slack", - channel_ref: str | None = "slack-ts-1", + channel_ref: str | None = None, ) -> None: conn = connect(db_path) + # Default to a per-question channel_ref so two OPEN rows never collide on the + # partial unique index uq_pending_questions_open_channel_ref (production never + # assigns the same channel_ref to two open questions — each gets its own + # posted message ts). Tests that need a specific ref still pass one. + if channel_ref is None: + channel_ref = f"slack-ts-{qid}" try: conn.execute( "INSERT INTO pending_questions " diff --git a/agent-team/tests/test_schema.py b/agent-team/tests/test_schema.py index 618e5fc..52c6e61 100644 --- a/agent-team/tests/test_schema.py +++ b/agent-team/tests/test_schema.py @@ -350,3 +350,116 @@ def test_shared_connection_concurrent_same_question_single_winner( assert errors == [], f"shared-connection CAS raised: {errors!r}" assert sum(wins) == 1 assert wins.count(False) == n - 1 + + +def _insert_open_question_with_ref( + conn: sqlite3.Connection, + qid: str, + channel_ref: str | None, + *, + status: str = "open", + turn: int = 0, +) -> None: + conn.execute( + "INSERT INTO pending_questions " + "(question_id, thread_id, turn, status, transport, channel_ref) " + "VALUES (?, 'thread-1', ?, ?, 'slack', ?)", + (qid, turn, status, channel_ref), + ) + + +def test_open_channel_ref_partial_unique_rejects_second_open_row( + tmp_path: Path, +) -> None: + """Defense-in-depth: two OPEN rows can never share a non-null channel_ref. + + The partial unique index ``uq_pending_questions_open_channel_ref`` makes a + duplicate (channel_ref, status='open') pair impossible, so a thread reply's + thread_ts can never resolve to two open questions. + """ + db = tmp_path / "db.sqlite" + init_db(db) + conn = connect(db) + try: + _insert_open_question_with_ref(conn, "q1", "ts-100") + with pytest.raises(sqlite3.IntegrityError): + _insert_open_question_with_ref(conn, "q2", "ts-100") + finally: + conn.close() + + +def test_open_channel_ref_partial_unique_allows_multiple_nulls( + tmp_path: Path, +) -> None: + """NULL channel_refs are unaffected: many open rows may have a null ref. + + The WHERE clause excludes nulls (and SQLite treats multiple NULLs as + distinct in a unique index anyway), so the lost-post recovery path's + unposted-open rows are never blocked. + """ + db = tmp_path / "db.sqlite" + init_db(db) + conn = connect(db) + try: + _insert_open_question_with_ref(conn, "q1", None) + _insert_open_question_with_ref(conn, "q2", None) # must not raise + count = conn.execute( + "SELECT COUNT(*) FROM pending_questions WHERE channel_ref IS NULL" + ).fetchone()[0] + finally: + conn.close() + assert count == 2 + + +def test_open_channel_ref_partial_unique_ignores_closed_rows( + tmp_path: Path, +) -> None: + """Closed rows are excluded: a non-open row may share an open row's ref. + + The index predicate is ``status='open'``, so once a question is answered / + expired / superseded its channel_ref no longer participates — a fresh open + question can reuse it (e.g. a re-asked turn on the same thread anchor). + """ + db = tmp_path / "db.sqlite" + init_db(db) + conn = connect(db) + try: + # An answered row and a superseded row both carry 'ts-200'. + _insert_open_question_with_ref(conn, "ans", "ts-200", status="answered") + _insert_open_question_with_ref(conn, "sup", "ts-200", status="superseded") + # A NEW open row may still take 'ts-200' (no open row holds it). + _insert_open_question_with_ref(conn, "open1", "ts-200") # must not raise + # But a SECOND open row with the same ref is rejected. + with pytest.raises(sqlite3.IntegrityError): + _insert_open_question_with_ref(conn, "open2", "ts-200") + finally: + conn.close() + + +def test_migrate_adds_open_channel_ref_partial_unique(tmp_path: Path) -> None: + """migrate() (not just init_db) installs the partial unique index. + + Existing DBs picked up via migrate() must get the defense-in-depth index too, + so the uniqueness guarantee holds after an in-place schema step. + """ + db = tmp_path / "db.sqlite" + conn = connect(db) + try: + # Migrate twice: the second call runs against an already-stamped v1 DB, + # exercising the UNCONDITIONAL index install (the `current < 1` block is + # skipped), which is exactly the existing-DB upgrade path. + migrate(conn) + migrate(conn) + names = { + r[0] + for r in conn.execute( + "SELECT name FROM sqlite_master WHERE type='index'" + ).fetchall() + } + assert "uq_pending_questions_open_channel_ref" in names + # And it actually enforces: a duplicate open ref is rejected. + _insert_open_question_with_ref(conn, "q1", "ts-300") + with pytest.raises(sqlite3.IntegrityError): + _insert_open_question_with_ref(conn, "q2", "ts-300") + finally: + conn.close() diff --git a/agent-team/tests/test_slack_listener.py b/agent-team/tests/test_slack_listener.py index 911bdc6..9b7cadc 100644 --- a/agent-team/tests/test_slack_listener.py +++ b/agent-team/tests/test_slack_listener.py @@ -622,3 +622,68 @@ def test_close_stops_socket_handler(db_path: Path) -> None: 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) + + +# --------------------------------------------------------------------------- +# handle_event — broadened guard (CWE-248): a racing submit never escapes. +# --------------------------------------------------------------------------- + + +def test_handle_event_swallows_keyerror_from_submit_answer( + db_path: Path, + monkeypatch: pytest.MonkeyPatch, + caplog: pytest.LogCaptureFixture, +) -> None: + """A KeyError from submit_answer (concurrent row mutation) is logged-and-ignored. + + On the accept path submit_answer -> _question_turn / _question_thread raise + KeyError if the just-answered row is concurrently deleted/mutated. Such an + exception must NOT escape handle_event into the Bolt dispatch (it would take + the listener down). The broadened guard returns None at WARNING (loud, not + silent), without widening the trust boundary — authorization already ran. + """ + import logging + + import agent_team.transport.slack_listener as slack_listener_module + + _seed_open_question(db_path, question_id="q1", thread_id="t1", turn=0) + queue = RecordingQueue() + listener = _listener(db_path, queue) + + def _boom(*_args: Any, **_kwargs: Any) -> Any: + raise KeyError("unknown question_id: 'q1'") + + monkeypatch.setattr(slack_listener_module, "submit_answer", _boom) + + with caplog.at_level(logging.WARNING): + outcome = listener.handle_event(_interactive_payload("q1", value="approve")) + + # No crash: the malformed/racing event is dropped, not raised. + assert outcome is None + assert queue.jobs == [] + # Surfaced at WARNING (not silent) so a real bug is still visible. + assert any( + r.levelno == logging.WARNING and "unexpected error" in r.message + for r in caplog.records + ) + + +def test_handle_event_swallows_sqlite_error_from_submit_answer( + db_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + """A sqlite3.Error from the underlying CAS is also caught, never escaping.""" + import sqlite3 + + import agent_team.transport.slack_listener as slack_listener_module + + _seed_open_question(db_path, question_id="q1", thread_id="t1", turn=0) + listener = _listener(db_path, RecordingQueue()) + + def _boom(*_args: Any, **_kwargs: Any) -> Any: + raise sqlite3.OperationalError("disk I/O error") + + monkeypatch.setattr(slack_listener_module, "submit_answer", _boom) + + # Must not raise; returns None. + assert listener.handle_event(_interactive_payload("q1")) is None -- 2.50.1