diff --git a/.env.example b/.env.example index 893bda621..e2f85a48b 100644 --- a/.env.example +++ b/.env.example @@ -361,6 +361,8 @@ BROWSER_INACTIVITY_TIMEOUT=120 # TELEGRAM_HOME_CHANNEL= # Default chat for cron delivery # TELEGRAM_HOME_CHANNEL_NAME= # Display name for home channel # TELEGRAM_CRON_THREAD_ID= # Forum topic ID for cron deliveries; overrides TELEGRAM_HOME_CHANNEL_THREAD_ID for cron so replies work in topic mode +# TELEGRAM_ALLOW_BOTS=none # Accept messages from other bots: none (default) | mentions | all +# TELEGRAM_GROUP_AUTOAPPROVE=false # When true, gate "bot added to group" through HSM auto-approval (requires HSM_URL + HERMES_AGENT_NAME); adder must be a platform admin or the bot declines and leaves. Default false = log only. # Webhook mode (optional — for cloud deployments like Fly.io/Railway) # Default is long polling. Setting TELEGRAM_WEBHOOK_URL switches to webhook mode. diff --git a/plugins/platforms/telegram/adapter.py b/plugins/platforms/telegram/adapter.py index f2b9800f2..a3afc443b 100644 --- a/plugins/platforms/telegram/adapter.py +++ b/plugins/platforms/telegram/adapter.py @@ -215,6 +215,7 @@ async def _shutdown_abandoned_app(app) -> None: Application, CommandHandler, CallbackQueryHandler, + ChatMemberHandler, MessageHandler as TelegramMessageHandler, ContextTypes, filters, @@ -233,6 +234,7 @@ async def _shutdown_abandoned_app(app) -> None: Application = Any CommandHandler = Any CallbackQueryHandler = Any + ChatMemberHandler = Any TelegramMessageHandler = Any HTTPXRequest = Any filters = None @@ -373,7 +375,8 @@ def check_telegram_requirements() -> bool: """ global TELEGRAM_AVAILABLE, Update, Bot, Message, InlineKeyboardButton global InlineKeyboardMarkup, LinkPreviewOptions, Application - global CommandHandler, CallbackQueryHandler, TelegramMessageHandler + global CommandHandler, CallbackQueryHandler, ChatMemberHandler + global TelegramMessageHandler global ContextTypes, filters, ParseMode, ChatType, HTTPXRequest if TELEGRAM_AVAILABLE: return True @@ -392,6 +395,7 @@ def check_telegram_requirements() -> bool: from telegram.ext import ( Application as _App, CommandHandler as _CH, CallbackQueryHandler as _CQH, + ChatMemberHandler as _CMH, MessageHandler as _MH, ContextTypes as _CT, filters as _filters, ) @@ -408,6 +412,7 @@ def check_telegram_requirements() -> bool: Application = _App CommandHandler = _CH CallbackQueryHandler = _CQH + ChatMemberHandler = _CMH TelegramMessageHandler = _MH ContextTypes = _CT filters = _filters @@ -3589,6 +3594,12 @@ def _with_limits(httpx_kwargs: Optional[dict] = None) -> dict: )) # Handle inline keyboard button callbacks (update prompts) self._app.add_handler(CallbackQueryHandler(self._handle_callback_query)) + # Bot's own membership changes (added to / removed from chats) — + # drives the HSM group auto-approval gate (TELEGRAM_GROUP_AUTOAPPROVE). + self._app.add_handler(ChatMemberHandler( + self._handle_my_chat_member, + ChatMemberHandler.MY_CHAT_MEMBER, + )) # Start polling — retry initialize() for transient TLS resets. # Each attempt is capped by _init_timeout so a single unreachable @@ -8127,6 +8138,115 @@ def _effective_update_message(self, update: Update) -> Optional[Message]: """ return getattr(update, "effective_message", None) or getattr(update, "message", None) + def _telegram_group_autoapprove_enabled(self) -> bool: + """Return whether bot-added-to-group events are gated through HSM.""" + return os.getenv("TELEGRAM_GROUP_AUTOAPPROVE", "false").lower() in {"true", "1", "yes", "on"} + + async def _handle_my_chat_member(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + """Handle changes to the bot's OWN chat membership (my_chat_member). + + When the bot is newly added to a group/supergroup and + ``TELEGRAM_GROUP_AUTOAPPROVE`` is enabled, ask HSM to auto-approve the + group: HSM checks whether the adder is a platform admin and, if so, + adds the group to the allowlist. On approval the bot greets the group; + otherwise it politely declines and leaves (fail-closed — any error on + the enforcement path counts as not approved). + + With the flag unset (default) the event is only logged, preserving + prior behavior for fleets that deploy this adapter before the HSM + endpoint is live. + """ + change = getattr(update, "my_chat_member", None) + if change is None: + return + + chat = getattr(change, "chat", None) + chat_type = str(getattr(chat, "type", "") or "").lower() + if chat_type not in ("group", "supergroup"): + return + + old_status = str( + getattr(getattr(change, "old_chat_member", None), "status", "") or "" + ).lower() + new_status = str( + getattr(getattr(change, "new_chat_member", None), "status", "") or "" + ).lower() + + # Only act on transitions INTO the chat. Promotions/demotions + # (member <-> administrator) re-fire my_chat_member and must not + # re-trigger approval; the bot leaving/being kicked is ignored too. + if new_status not in ("member", "administrator"): + return + if old_status not in ("", "left", "kicked"): + return + + chat_id = getattr(chat, "id", None) + if chat_id is None: + return + group_id = str(chat_id) + adder = str(getattr(getattr(change, "from_user", None), "id", "") or "") + + if not self._telegram_group_autoapprove_enabled(): + logger.info( + "[%s] Bot added to Telegram group %s by user %s — " + "auto-approval disabled (TELEGRAM_GROUP_AUTOAPPROVE unset), no action", + self.name, group_id, adder or "unknown", + ) + return + + approved = False + try: + from plugins.swarm_map_policy import approve_group_add + approved = approve_group_add(group_id, adder, platform="telegram") + except Exception as e: + logger.warning( + "[%s] Group auto-approval check failed for %s (fail-closed): %s", + self.name, group_id, e, + ) + approved = False + + if approved: + logger.info( + "[%s] Telegram group %s approved via HSM (added by %s) — staying", + self.name, group_id, adder or "unknown", + ) + try: + await self._bot.send_message( + chat_id=chat_id, + text="Hi! This group has been approved for me — I'm ready to help.", + ) + except Exception as e: + logger.warning( + "[%s] Failed to send group greeting to %s: %s", + self.name, group_id, _redact_telegram_error_text(e), + ) + return + + logger.info( + "[%s] Telegram group %s not approved via HSM (added by %s) — leaving", + self.name, group_id, adder or "unknown", + ) + try: + await self._bot.send_message( + chat_id=chat_id, + text=( + "Sorry — this group isn't approved for me yet. " + "Ask my admin to approve it, then add me again." + ), + ) + except Exception as e: + logger.warning( + "[%s] Failed to send decline notice to %s: %s", + self.name, group_id, _redact_telegram_error_text(e), + ) + try: + await self._bot.leave_chat(chat_id=chat_id) + except Exception as e: + logger.warning( + "[%s] Failed to leave unapproved group %s: %s", + self.name, group_id, _redact_telegram_error_text(e), + ) + async def _handle_text_message(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: """Handle incoming text messages. diff --git a/plugins/swarm_map_policy/__init__.py b/plugins/swarm_map_policy/__init__.py index 0dc237d15..6dc17eec5 100644 --- a/plugins/swarm_map_policy/__init__.py +++ b/plugins/swarm_map_policy/__init__.py @@ -83,6 +83,57 @@ def is_group_allowed(group_id: str, platform: str) -> bool: return False +def approve_group_add( + group_id: str, added_by_user_id: str, platform: str = "telegram" +) -> bool: + """Request HSM auto-approval for a group the bot was just added to. + + Called when someone adds the bot to a new group. HSM verifies that the + adder is a platform admin and, if so, adds the group to the allowlist. + Fail-closed: returns True only on HTTP 200 with ``approved: true`` — + network errors, non-200 responses, and missing fields all deny. + """ + url = _hsm_url() + harness = _harness_id() + if not url or not harness: + logger.warning("swarm-map-policy: HSM not configured, denying group add") + return False + try: + resp = requests.post( + f"{url}/api/harnesses/{harness}/surfaces/{platform}/groups/{group_id}", + json={"addedByUserId": added_by_user_id}, + timeout=5, + ) + if resp.status_code != 200: + reason = "" + try: + reason = resp.json().get("error", "") + except Exception: + pass + logger.info( + "swarm-map-policy: group add denied for %s (HTTP %s%s)", + group_id, resp.status_code, f": {reason}" if reason else "", + ) + return False + data = resp.json() + if data.get("approved") is True: + logger.info( + "swarm-map-policy: group add approved for %s (already_allowed=%s restarted=%s)", + group_id, data.get("already_allowed", False), data.get("restarted"), + ) + return True + logger.info( + "swarm-map-policy: group add not approved for %s: %s", + group_id, data.get("reason", "no reason given"), + ) + return False + except Exception as e: + logger.warning( + "swarm-map-policy: group add approval failed (fail-closed): %s", e + ) + return False + + def is_tool_allowed(tool_name: str, group_id: str) -> bool: """Check if a tool is allowed for a group. Fail-open.""" url = _hsm_url() diff --git a/tests/gateway/test_telegram_group_autoapprove.py b/tests/gateway/test_telegram_group_autoapprove.py new file mode 100644 index 000000000..a6ce83438 --- /dev/null +++ b/tests/gateway/test_telegram_group_autoapprove.py @@ -0,0 +1,189 @@ +"""Tests for Telegram group auto-approval on bot add (my_chat_member). + +When an admin adds the bot to a group and TELEGRAM_GROUP_AUTOAPPROVE is +enabled, the adapter asks HSM to approve the group (the adder must be a +platform admin). Approved -> greet and stay; denied -> polite decline and +leave. Flag unset -> log-only no-op, preserving prior behavior. +""" +import os +from types import SimpleNamespace +from unittest.mock import AsyncMock, patch + +import pytest + +from gateway.config import Platform, PlatformConfig + + +def _make_adapter(**extra): + from plugins.platforms.telegram.adapter import TelegramAdapter + + adapter = object.__new__(TelegramAdapter) + adapter.platform = Platform.TELEGRAM + adapter.config = PlatformConfig(enabled=True, token="fake-token", extra=extra) + adapter._bot = SimpleNamespace( + id=999, + username="test_bot", + send_message=AsyncMock(), + leave_chat=AsyncMock(), + ) + return adapter + + +def _make_update( + *, + chat_id=-100123, + chat_type="supergroup", + old_status="left", + new_status="member", + adder_id=777, +): + return SimpleNamespace( + my_chat_member=SimpleNamespace( + chat=SimpleNamespace(id=chat_id, type=chat_type, title="Test Group"), + from_user=SimpleNamespace(id=adder_id), + old_chat_member=SimpleNamespace(status=old_status), + new_chat_member=SimpleNamespace(status=new_status), + ) + ) + + +_FLAG_ON = {"TELEGRAM_GROUP_AUTOAPPROVE": "1"} + + +@pytest.mark.asyncio +async def test_added_and_approved_greets_and_stays(): + """Approved add: bot greets the group and does not leave.""" + adapter = _make_adapter() + update = _make_update() + with patch.dict(os.environ, _FLAG_ON): + with patch( + "plugins.swarm_map_policy.approve_group_add", return_value=True + ) as mock_approve: + await adapter._handle_my_chat_member(update, None) + + mock_approve.assert_called_once_with("-100123", "777", platform="telegram") + adapter._bot.send_message.assert_awaited_once() + assert adapter._bot.send_message.await_args.kwargs["chat_id"] == -100123 + adapter._bot.leave_chat.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_added_and_denied_declines_and_leaves(): + """Denied add: bot sends one polite decline, then leaves the chat.""" + adapter = _make_adapter() + update = _make_update() + with patch.dict(os.environ, _FLAG_ON): + with patch("plugins.swarm_map_policy.approve_group_add", return_value=False): + await adapter._handle_my_chat_member(update, None) + + adapter._bot.send_message.assert_awaited_once() + text = adapter._bot.send_message.await_args.kwargs["text"] + assert "approved" in text.lower() + adapter._bot.leave_chat.assert_awaited_once_with(chat_id=-100123) + + +@pytest.mark.asyncio +async def test_flag_off_is_logged_noop(): + """Flag unset: event is a no-op — no HSM call, no message, no leave.""" + adapter = _make_adapter() + update = _make_update() + env = {k: v for k, v in os.environ.items() if k != "TELEGRAM_GROUP_AUTOAPPROVE"} + with patch.dict(os.environ, env, clear=True): + with patch( + "plugins.swarm_map_policy.approve_group_add", return_value=False + ) as mock_approve: + await adapter._handle_my_chat_member(update, None) + + mock_approve.assert_not_called() + adapter._bot.send_message.assert_not_awaited() + adapter._bot.leave_chat.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_promotion_transition_is_noop(): + """member -> administrator (promotion) must not re-trigger approval.""" + adapter = _make_adapter() + update = _make_update(old_status="member", new_status="administrator") + with patch.dict(os.environ, _FLAG_ON): + with patch( + "plugins.swarm_map_policy.approve_group_add", return_value=False + ) as mock_approve: + await adapter._handle_my_chat_member(update, None) + + mock_approve.assert_not_called() + adapter._bot.send_message.assert_not_awaited() + adapter._bot.leave_chat.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_bot_leaving_is_noop(): + """member -> left (bot removed) must not trigger anything.""" + adapter = _make_adapter() + update = _make_update(old_status="member", new_status="left") + with patch.dict(os.environ, _FLAG_ON): + with patch( + "plugins.swarm_map_policy.approve_group_add", return_value=True + ) as mock_approve: + await adapter._handle_my_chat_member(update, None) + + mock_approve.assert_not_called() + adapter._bot.send_message.assert_not_awaited() + adapter._bot.leave_chat.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_private_chat_is_noop(): + """my_chat_member in a private chat (user unblocked bot etc.) is ignored.""" + adapter = _make_adapter() + update = _make_update(chat_type="private", chat_id=555) + with patch.dict(os.environ, _FLAG_ON): + with patch( + "plugins.swarm_map_policy.approve_group_add", return_value=True + ) as mock_approve: + await adapter._handle_my_chat_member(update, None) + + mock_approve.assert_not_called() + adapter._bot.send_message.assert_not_awaited() + adapter._bot.leave_chat.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_plain_group_chat_is_enforced(): + """Basic (non-super) groups are gated too.""" + adapter = _make_adapter() + update = _make_update(chat_type="group", chat_id=-4567) + with patch.dict(os.environ, _FLAG_ON): + with patch( + "plugins.swarm_map_policy.approve_group_add", return_value=True + ) as mock_approve: + await adapter._handle_my_chat_member(update, None) + + mock_approve.assert_called_once_with("-4567", "777", platform="telegram") + + +@pytest.mark.asyncio +async def test_approval_error_fails_closed(): + """Exception in the approval path counts as denied: decline + leave.""" + adapter = _make_adapter() + update = _make_update() + with patch.dict(os.environ, _FLAG_ON): + with patch( + "plugins.swarm_map_policy.approve_group_add", + side_effect=Exception("HSM exploded"), + ): + await adapter._handle_my_chat_member(update, None) + + adapter._bot.leave_chat.assert_awaited_once_with(chat_id=-100123) + + +@pytest.mark.asyncio +async def test_decline_send_failure_still_leaves(): + """If the decline message can't be sent, the bot still leaves.""" + adapter = _make_adapter() + adapter._bot.send_message.side_effect = Exception("send failed") + update = _make_update() + with patch.dict(os.environ, _FLAG_ON): + with patch("plugins.swarm_map_policy.approve_group_add", return_value=False): + await adapter._handle_my_chat_member(update, None) + + adapter._bot.leave_chat.assert_awaited_once_with(chat_id=-100123) diff --git a/tests/plugins/test_swarm_map_policy.py b/tests/plugins/test_swarm_map_policy.py index 4443e22bb..d01819087 100644 --- a/tests/plugins/test_swarm_map_policy.py +++ b/tests/plugins/test_swarm_map_policy.py @@ -65,6 +65,86 @@ def test_admin_check_fail_closed(self): assert is_platform_admin("user-123", "signal") is False +class TestApproveGroupAdd: + """Tests for approve_group_add — HSM group auto-approval on bot add.""" + + def _call(self, mock_resp=None, side_effect=None, group_id="-100123", adder="777"): + from plugins.swarm_map_policy import approve_group_add + with patch("plugins.swarm_map_policy._hsm_url", return_value="http://hsm:3002"), \ + patch("plugins.swarm_map_policy._harness_id", return_value="hermes-test"), \ + patch("plugins.swarm_map_policy.requests") as mock_req: + if side_effect is not None: + mock_req.post.side_effect = side_effect + else: + mock_req.post.return_value = mock_resp + result = approve_group_add(group_id, adder) + return result, mock_req + + @staticmethod + def _resp(status_code=200, body=None): + resp = MagicMock() + resp.status_code = status_code + resp.json.return_value = body if body is not None else {} + return resp + + def test_approved(self): + """200 + approved:true returns True.""" + result, _ = self._call(self._resp(200, {"approved": True, "restarted": True})) + assert result is True + + def test_approved_already_allowed(self): + """200 + approved:true + already_allowed (wildcard/listed) returns True.""" + result, _ = self._call(self._resp(200, {"approved": True, "already_allowed": True})) + assert result is True + + def test_approved_restart_failed_still_true(self): + """approved:true with restarted:false (env written, recreate failed) is still approved.""" + result, _ = self._call(self._resp(200, {"approved": True, "restarted": False})) + assert result is True + + def test_not_approved(self): + """200 + approved:false returns False.""" + result, _ = self._call( + self._resp(200, {"approved": False, "reason": "adder is not an admin"}) + ) + assert result is False + + def test_bad_request_denied(self): + """400 with error body is treated as not approved.""" + result, _ = self._call(self._resp(400, {"error": "missing addedByUserId"})) + assert result is False + + def test_network_error_fail_closed(self): + """Network failure denies (fail-closed).""" + result, _ = self._call(side_effect=Exception("Connection refused")) + assert result is False + + def test_missing_approved_field_denied(self): + """200 with no approved field denies (fail-closed).""" + result, _ = self._call(self._resp(200, {})) + assert result is False + + def test_non_bool_approved_denied(self): + """approved must be strictly true — truthy strings deny.""" + result, _ = self._call(self._resp(200, {"approved": "yes"})) + assert result is False + + def test_no_config_fail_closed(self): + """Missing HSM_URL denies without any request.""" + from plugins.swarm_map_policy import approve_group_add + with patch("plugins.swarm_map_policy._hsm_url", return_value=None): + assert approve_group_add("-100123", "777") is False + + def test_posts_correct_url_and_body(self): + """Request hits the HSM group endpoint with addedByUserId body.""" + _, mock_req = self._call(self._resp(200, {"approved": True})) + mock_req.post.assert_called_once_with( + "http://hsm:3002/api/harnesses/hermes-test/surfaces/telegram/groups/-100123", + json={"addedByUserId": "777"}, + timeout=5, + ) + + class TestSessionContextCaching: """Tests for session context caching via pre_gateway_dispatch."""