fix(agent-team): harden listener respawn/close, broaden handle_event guard, channel_ref partial-unique (#29)
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).
This commit is contained in:
parent
8afe876fa2
commit
a6275b4000
8 changed files with 352 additions and 3 deletions
|
|
@ -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).
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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 "
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Reference in a new issue