fix(gateway): add bot-to-bot loop guard for Telegram
Telegram Bot API 10.0 (2026-05-08) lets bots receive messages from other
bots, and the platform ships no loop guard of its own. core.telegram.org
/api/bots/bot-to-bot ("Loop prevention") requires the BOT to make
bot-message handling terminate predictably via dedupe, per-chat rate
limits and maximum interaction depth, and warns that "failure to handle
loops properly may lead to degraded performance or platform
restrictions".
Hermes had no brake at all: TELEGRAM_ALLOW_BOTS in gateway/authz_mixin.py
was a bare `return True` with no counter, window or cooldown.
TELEGRAM_ALLOW_BOTS=mentions does not help, because when bot A replies to
bot B the reply itself satisfies the mention test, so every turn re-arms
the peer. Observed 2026-08-21: two Hermes bots exchanged 132 messages in
one group before a human intervened.
This adds gateway/bot_loop_guard.py, a thread-safe sliding-window budget
with cooldown and sweeping of stale buckets, wired into
_is_user_authorized. Defaults (20 events / 60 s window / 60 s cooldown)
match the OpenClaw reference implementation.
Two placement details that are easy to get wrong:
- The guard runs BEFORE the TELEGRAM_GROUP_ALLOWED_CHATS shortcut. That
env var returns True for any sender in an allowlisted chat, bots
included, ~32 lines before the is_bot block, so a guard placed in the
is_bot block never executes for the one configuration that needs it.
- The budget is keyed per CONVERSATION, not per (sender, receiver). This
process only sees inbound messages, so the receiver is constant; keying
on the pair would give each sender its own budget and N bots would need
N x budget messages to trip a guard meant to cap the whole exchange.
Installations without ALLOW_BOTS are byte-identical in behaviour: the
guard only runs on the source.is_bot branch with ALLOW_BOTS set to
mentions or all. HERMES_BOT_LOOP_PROTECTION=off is a full kill switch.
Adds tests/gateway/test_bot_loop_guard.py (7 cases), each verified
failing against the tree without this patch.
This commit is contained in:
@@ -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
|
||||
|
||||
|
||||
@@ -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()
|
||||
@@ -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"
|
||||
Reference in New Issue
Block a user