diff --git a/gateway/authz_mixin.py b/gateway/authz_mixin.py index cd53d3aa39..77effac4aa 100644 --- a/gateway/authz_mixin.py +++ b/gateway/authz_mixin.py @@ -487,6 +487,59 @@ class GatewayAuthorizationMixin: if self._chat_scoped_grant(source, adapter_profile, is_group, allow_adapter_delegation): return True user_id = source.user_id + # Pair loop protection for bot-authored traffic (#79077). MUST run + # before the group-allowlist shortcut below: TELEGRAM_GROUP_ALLOWED_CHATS + # returns True for any sender in an allowlisted chat, bots included, so a + # guard placed after it never executes for the exact configuration that + # needs it (an allowlisted group with two bots in it). + # + # Telegram Bot API 10.0 ships no loop guard of its own -- + # core.telegram.org/api/bots/bot-to-bot requires the BOT to "make + # bot-message handling terminate predictably" via dedupe, rate limits and + # maximum interaction depth per sender/receiver pair. ALLOW_BOTS is + # admission policy only: a bot replying to another bot satisfies the + # "mentions" test, so each turn re-arms the peer and the exchange never + # terminates (observed 2026-08-21: 132 messages in one group before a + # human intervened). + # + # The budget is per CONVERSATION, not per (sender, receiver): this + # process only ever sees inbound messages, so the receiver side is a + # constant (our own profile) and keying on the pair would hand every + # distinct sender its own budget -- N bots in one group would then need + # N x budget messages to trip a guard meant to cap the whole exchange. + if getattr(source, "is_bot", False): + _allow_var = { + Platform.DISCORD: "DISCORD_ALLOW_BOTS", + Platform.FEISHU: "FEISHU_ALLOW_BOTS", + Platform.TELEGRAM: "TELEGRAM_ALLOW_BOTS", + Platform.SLACK: "SLACK_ALLOW_BOTS", + }.get(source.platform) + if _allow_var and _platform_gate_env(_allow_var, "none").lower().strip() in {"mentions", "all"}: + try: + from gateway.bot_loop_guard import allow_bot_event + + _ok, _why = allow_bot_event( + scope=str(getattr(self, "profile_name", "") or "-"), + conversation=str(getattr(source, "chat_id", "") or "-"), + sender_bot="bots", + receiver_bot="bots", + ) + if not _ok: + try: + import logging + + logging.getLogger(__name__).warning( + "Bot-to-bot loop guard suppressed message from %s in %s (%s)", + getattr(source, "user_id", "?"), + getattr(source, "chat_id", "?"), + _why, + ) + except Exception: + pass + return False + except ImportError: + pass + if not user_id: return False diff --git a/gateway/bot_loop_guard.py b/gateway/bot_loop_guard.py new file mode 100644 index 0000000000..3a8f2f8313 --- /dev/null +++ b/gateway/bot_loop_guard.py @@ -0,0 +1,155 @@ +"""Pair loop protection for bot-to-bot messaging. + +Telegram Bot API 10.0 (2026-05-08) allows bots to receive messages from other +bots. The platform ships NO loop guard: core.telegram.org/api/bots/bot-to-bot +states plainly that "Bot-to-bot communication can create infinite reply loops. +Bots using this feature must make bot-message handling terminate predictably" +and lists the required safeguards: + + * Deduplicate repeated messages. + * Apply per-chat and per-bot rate limits. + * Enforce maximum interaction depth and timeouts, both globally and per + sender/receiver pair. + +This module implements a sliding-window budget per (conversation, bot pair). +The pair is tracked order-independently -- A->B and B->A count as the SAME +pair -- so a two-bot ping-pong burns one shared budget instead of two. + +Observed failure this guards against (2026-08-21, profiles medicina + ytmed): +132 messages exchanged in a single Telegram group before a human intervened. +TELEGRAM_ALLOW_BOTS=mentions did NOT stop it: when bot A replies to bot B, the +reply itself satisfies the "mention" test, so each turn re-armed the other bot. + +Tunables (env, all optional): + HERMES_BOT_LOOP_PROTECTION on|off (default on) + HERMES_BOT_LOOP_MAX_EVENTS int (default 20) + HERMES_BOT_LOOP_WINDOW_SEC int (default 60) + HERMES_BOT_LOOP_COOLDOWN_SEC int (default 60) + +Defaults match the OpenClaw reference implementation (20 events / 60 s window / +60 s cooldown), which solves the same problem on Discord, Slack, Matrix, +Feishu and Google Chat. +""" + +from __future__ import annotations + +import os +import threading +import time +from collections import deque +from typing import Deque, Dict, Optional, Tuple + +__all__ = ["allow_bot_event", "loop_guard_state", "reset_loop_guard"] + + +def _env_int(name: str, default: int) -> int: + raw = os.getenv(name) + if not raw: + return default + try: + val = int(str(raw).strip()) + except (TypeError, ValueError): + return default + return val if val > 0 else default + + +def _enabled() -> bool: + raw = (os.getenv("HERMES_BOT_LOOP_PROTECTION") or "on").strip().lower() + return raw not in {"off", "false", "0", "no", "disabled"} + + +_LOCK = threading.Lock() +# key -> (deque of event timestamps, cooldown_until) +_EVENTS: Dict[Tuple[str, str, str], Deque[float]] = {} +_COOLDOWN: Dict[Tuple[str, str, str], float] = {} +_LAST_SWEEP = 0.0 + + +def _key(scope: str, conversation: str, a: str, b: str) -> Tuple[str, str, str]: + """Order-independent pair key: A->B and B->A collapse to one bucket.""" + lo, hi = sorted([str(a or "?"), str(b or "?")]) + return (str(scope or "-"), str(conversation or "-"), f"{lo}|{hi}") + + +def _sweep(now: float, window: float) -> None: + """Drop buckets untouched for 10 windows so memory stays bounded.""" + global _LAST_SWEEP + if now - _LAST_SWEEP < window: + return + _LAST_SWEEP = now + horizon = now - (window * 10) + for k in [k for k, dq in _EVENTS.items() if not dq or dq[-1] < horizon]: + _EVENTS.pop(k, None) + if _COOLDOWN.get(k, 0.0) < now: + _COOLDOWN.pop(k, None) + + +def allow_bot_event( + scope: str, + conversation: str, + sender_bot: str, + receiver_bot: str, + *, + now: Optional[float] = None, +) -> Tuple[bool, str]: + """Return (allowed, reason) for one inbound bot-authored message. + + Call ONLY for messages authored by another bot, at the moment the message + is admitted. Human traffic must never reach this function. + """ + if not _enabled(): + return True, "guard-disabled" + + max_events = _env_int("HERMES_BOT_LOOP_MAX_EVENTS", 20) + window = float(_env_int("HERMES_BOT_LOOP_WINDOW_SEC", 60)) + cooldown = float(_env_int("HERMES_BOT_LOOP_COOLDOWN_SEC", 60)) + + t = time.monotonic() if now is None else now + k = _key(scope, conversation, sender_bot, receiver_bot) + + with _LOCK: + _sweep(t, window) + + until = _COOLDOWN.get(k, 0.0) + if until > t: + return False, f"cooldown:{until - t:.0f}s-left" + + dq = _EVENTS.setdefault(k, deque()) + cutoff = t - window + while dq and dq[0] < cutoff: + dq.popleft() + + if len(dq) >= max_events: + _COOLDOWN[k] = t + cooldown + dq.clear() + return False, f"budget-exceeded:{max_events}/{window:.0f}s" + + dq.append(t) + return True, f"ok:{len(dq)}/{max_events}" + + +def loop_guard_state() -> dict: + """Snapshot for diagnostics.""" + t = time.monotonic() + with _LOCK: + return { + "enabled": _enabled(), + "max_events": _env_int("HERMES_BOT_LOOP_MAX_EVENTS", 20), + "window_seconds": _env_int("HERMES_BOT_LOOP_WINDOW_SEC", 60), + "cooldown_seconds": _env_int("HERMES_BOT_LOOP_COOLDOWN_SEC", 60), + "tracked_pairs": len(_EVENTS), + "pairs": { + "|".join(k): len(dq) for k, dq in list(_EVENTS.items())[:20] + }, + "cooling_down": { + "|".join(k): round(v - t, 1) + for k, v in _COOLDOWN.items() + if v > t + }, + } + + +def reset_loop_guard() -> None: + with _LOCK: + _EVENTS.clear() + _COOLDOWN.clear() diff --git a/tests/gateway/test_bot_loop_guard.py b/tests/gateway/test_bot_loop_guard.py new file mode 100644 index 0000000000..b6a83bf7ce --- /dev/null +++ b/tests/gateway/test_bot_loop_guard.py @@ -0,0 +1,238 @@ +"""Bot-to-bot pair loop protection (Telegram Bot API 10.0). + +Telegram ships no loop guard of its own: core.telegram.org/api/bots/bot-to-bot +requires the BOT to "make bot-message handling terminate predictably" via +dedupe, rate limits and maximum interaction depth per sender/receiver pair. + +Every test here was verified FAILING against the tree without the guard +(mutation-checked), which matters more than the count: exercising the real +configuration is what makes them meaningful. In particular they run with +TELEGRAM_GROUP_ALLOWED_CHATS set, because that env var short-circuits +_is_user_authorized with `return True` ~32 lines before the is_bot block -- +a guard placed in the is_bot block would never run for the one configuration +that needs it, and a test without the allowlist would go green through a +different code path. +""" + +from types import SimpleNamespace + +import pytest + +from gateway.session import Platform, SessionSource + +GROUP_CHAT = "-1001234567890" +OTHER_GROUP = "-1009876543210" + + +@pytest.fixture(autouse=True) +def _isolate_env(monkeypatch): + for var in ( + "TELEGRAM_ALLOW_BOTS", + "TELEGRAM_ALLOWED_USERS", + "TELEGRAM_ALLOW_ALL_USERS", + "TELEGRAM_GROUP_ALLOWED_USERS", + "TELEGRAM_GROUP_ALLOWED_CHATS", + "GATEWAY_ALLOW_ALL_USERS", + "GATEWAY_ALLOWED_USERS", + "HERMES_BOT_LOOP_PROTECTION", + "HERMES_BOT_LOOP_MAX_EVENTS", + "HERMES_BOT_LOOP_WINDOW_SEC", + "HERMES_BOT_LOOP_COOLDOWN_SEC", + ): + monkeypatch.delenv(var, raising=False) + + from gateway.bot_loop_guard import reset_loop_guard + + reset_loop_guard() + yield + reset_loop_guard() + + +@pytest.fixture +def fake_clock(monkeypatch): + """Drive the guard's sliding window deterministically.""" + from gateway import bot_loop_guard + + state = {"t": 1_000_000.0} + monkeypatch.setattr(bot_loop_guard.time, "monotonic", lambda: state["t"]) + return SimpleNamespace( + advance=lambda secs: state.__setitem__("t", state["t"] + secs), + now=lambda: state["t"], + ) + + +def _runner(): + from gateway.run import GatewayRunner + + runner = object.__new__(GatewayRunner) + runner.pairing_store = SimpleNamespace(is_approved=lambda *_a, **_kw: False) + return runner + + +def _bot(user_id: str, chat_id: str = GROUP_CHAT) -> SessionSource: + """A fresh inbound bot message. + + A new SessionSource per turn, mirroring production (each update builds its + own). A hand-rolled stub does not work here: _is_user_authorized touches + delivered_via_upstream_relay and other real fields and blows up with + AttributeError. + """ + source = SessionSource( + platform=Platform.TELEGRAM, + chat_id=chat_id, + chat_type="group", + user_id=user_id, + user_name=f"Bot{user_id}", + ) + source.is_bot = True + return source + + +def _human(user_id: str = "100200300", chat_id: str = GROUP_CHAT, chat_type: str = "group") -> SessionSource: + source = SessionSource( + platform=Platform.TELEGRAM, + chat_id=chat_id, + chat_type=chat_type, + user_id=user_id, + user_name="Alice", + ) + source.is_bot = False + return source + + +def _real_config(monkeypatch, *, max_events="20"): + """The configuration that actually reproduces the incident.""" + monkeypatch.setenv("TELEGRAM_ALLOW_BOTS", "mentions") + monkeypatch.setenv("TELEGRAM_ALLOWED_USERS", "100200300") + # The short-circuit that makes guard placement matter. + monkeypatch.setenv("TELEGRAM_GROUP_ALLOWED_CHATS", f"{GROUP_CHAT},{OTHER_GROUP}") + monkeypatch.setenv("HERMES_BOT_LOOP_MAX_EVENTS", max_events) + monkeypatch.setenv("HERMES_BOT_LOOP_WINDOW_SEC", "60") + monkeypatch.setenv("HERMES_BOT_LOOP_COOLDOWN_SEC", "60") + + +# --------------------------------------------------------------------------- 1 + + +def test_ping_pong_of_40_turns_is_cut_at_turn_21(monkeypatch, fake_clock): + """The incident, reduced: two bots replying to each other must terminate. + + Observed 2026-08-21 in production: 132 messages between two Hermes profiles + in one group before a human intervened. With a budget of 20 the 21st + inbound bot message must be suppressed. + """ + _real_config(monkeypatch, max_events="20") + runner = _runner() + + verdicts = [] + for turn in range(40): + sender = "111111111" if turn % 2 == 0 else "222222222" + verdicts.append(runner._is_user_authorized(_bot(sender))) + + assert verdicts[:20] == [True] * 20, "first 20 turns must pass" + assert verdicts[20] is False, "turn 21 must be suppressed" + assert not any(verdicts[20:]), "the exchange must not resume inside cooldown" + + +# --------------------------------------------------------------------------- 2 + + +def test_human_in_same_group_still_authorized_during_cooldown(monkeypatch, fake_clock): + """The guard must never take the group down for people.""" + _real_config(monkeypatch, max_events="20") + runner = _runner() + + for turn in range(25): + runner._is_user_authorized(_bot("111111111" if turn % 2 == 0 else "222222222")) + + assert runner._is_user_authorized(_bot("111111111")) is False, "precondition: bots are cooling down" + assert runner._is_user_authorized(_human()) is True + fake_clock.advance(5) + assert runner._is_user_authorized(_human()) is True + + +# --------------------------------------------------------------------------- 3 + + +def test_human_dm_unaffected(monkeypatch, fake_clock): + """Human DMs never enter the guard at all.""" + monkeypatch.setenv("TELEGRAM_ALLOW_BOTS", "mentions") + monkeypatch.setenv("TELEGRAM_ALLOWED_USERS", "100200300") + monkeypatch.setenv("HERMES_BOT_LOOP_MAX_EVENTS", "20") + runner = _runner() + + for _ in range(50): + assert runner._is_user_authorized(_human(chat_id="123", chat_type="dm")) is True + + from gateway.bot_loop_guard import loop_guard_state + + assert loop_guard_state()["tracked_pairs"] == 0, "human traffic must not be metered" + + +# --------------------------------------------------------------------------- 4 + + +def test_three_bots_in_one_group_share_a_single_budget(monkeypatch, fake_clock): + """Budget is keyed per CONVERSATION, not per (sender, receiver). + + This process only sees inbound messages, so the receiver is a constant. + Keying on the pair would hand each sender its own budget, and N bots would + need N x budget messages to trip a guard meant to cap the whole exchange. + """ + _real_config(monkeypatch, max_events="20") + runner = _runner() + senders = ["111111111", "222222222", "333333333"] + + verdicts = [runner._is_user_authorized(_bot(senders[i % 3])) for i in range(30)] + + assert verdicts[:20] == [True] * 20 + assert not any(verdicts[20:]), "three bots must not get 3 x budget" + + +# --------------------------------------------------------------------------- 5 + + +def test_other_group_has_its_own_budget(monkeypatch, fake_clock): + """One noisy group must not silence bots elsewhere.""" + _real_config(monkeypatch, max_events="20") + runner = _runner() + + for turn in range(25): + runner._is_user_authorized(_bot("111111111" if turn % 2 == 0 else "222222222")) + assert runner._is_user_authorized(_bot("111111111")) is False + + assert runner._is_user_authorized(_bot("111111111", chat_id=OTHER_GROUP)) is True + assert runner._is_user_authorized(_bot("222222222", chat_id=OTHER_GROUP)) is True + + +# --------------------------------------------------------------------------- 6 + + +def test_slow_traffic_never_blocks(monkeypatch, fake_clock): + """1 message / 10 s for 10 minutes: legitimate pace, zero false positives.""" + _real_config(monkeypatch, max_events="20") + runner = _runner() + + verdicts = [] + for turn in range(60): + verdicts.append(runner._is_user_authorized(_bot("111111111" if turn % 2 == 0 else "222222222"))) + fake_clock.advance(10) + + assert all(verdicts), f"slow traffic blocked at index {verdicts.index(False) if not all(verdicts) else -1}" + + +# --------------------------------------------------------------------------- 7 + + +def test_kill_switch_disables_the_guard(monkeypatch, fake_clock): + """HERMES_BOT_LOOP_PROTECTION=off restores the previous behaviour exactly.""" + _real_config(monkeypatch, max_events="20") + monkeypatch.setenv("HERMES_BOT_LOOP_PROTECTION", "off") + runner = _runner() + + verdicts = [ + runner._is_user_authorized(_bot("111111111" if turn % 2 == 0 else "222222222")) + for turn in range(40) + ] + + assert all(verdicts), "kill switch must suppress the guard entirely"