From 56f31c0aa29513db58c3c3cd4605e3fba6fb8f0b Mon Sep 17 00:00:00 2001 From: "2569718930@qq.com" <2569718930@qq.com> Date: Sun, 14 Jun 2026 06:32:58 +0800 Subject: [PATCH] Add Telegram webhook command entrypoint --- src/bot/orchestrator.py | 18 ++++++ tests/test_telegram_webhook.py | 90 ++++++++++++++++++++++++++ web/app_factory.py | 2 + web/routers/telegram_webhook.py | 12 ++++ web/services/telegram_webhook.py | 108 +++++++++++++++++++++++++++++++ 5 files changed, 230 insertions(+) create mode 100644 tests/test_telegram_webhook.py create mode 100644 web/routers/telegram_webhook.py create mode 100644 web/services/telegram_webhook.py diff --git a/src/bot/orchestrator.py b/src/bot/orchestrator.py index fd2a1b6c..82cf9ac1 100644 --- a/src/bot/orchestrator.py +++ b/src/bot/orchestrator.py @@ -13,6 +13,20 @@ from src.utils.config_validation import validate_or_raise from src.utils.telegram_chat_ids import get_telegram_chat_ids_from_env +def _env_bool(name: str, default: bool) -> bool: + raw = os.getenv(name) + if raw is None: + return default + return raw.strip().lower() in {"1", "true", "yes", "on"} + + +def _wait_forever() -> None: + import time + + while True: + time.sleep(3600) + + def _register_handlers( bot: Any, config: dict[str, Any], @@ -73,4 +87,8 @@ def start_bot() -> None: started_count, len(runtime_status.loops), ) + if not _env_bool("TELEGRAM_BOT_POLLING_ENABLED", True): + logger.warning("Telegram bot polling disabled; background loops remain active") + _wait_forever() + return bot.infinity_polling(allowed_updates=["message", "callback_query", "chat_join_request"]) diff --git a/tests/test_telegram_webhook.py b/tests/test_telegram_webhook.py new file mode 100644 index 00000000..84b73bf7 --- /dev/null +++ b/tests/test_telegram_webhook.py @@ -0,0 +1,90 @@ +import sys +from types import SimpleNamespace + +from fastapi.testclient import TestClient + +from web.app import app +import web.services.telegram_webhook as telegram_webhook + + +client = TestClient(app) + + +def test_telegram_webhook_dispatches_update_when_secret_matches(monkeypatch): + calls = [] + + class FakeBot: + def process_new_updates(self, updates): + calls.append(updates) + + monkeypatch.setenv("TELEGRAM_WEBHOOK_SECRET", "secret-path") + monkeypatch.delenv("TELEGRAM_WEBHOOK_HEADER_SECRET", raising=False) + monkeypatch.setattr(telegram_webhook, "_get_webhook_bot", lambda: FakeBot()) + monkeypatch.setattr( + telegram_webhook, + "_parse_update", + lambda payload: {"update_id": payload["update_id"]}, + ) + + response = client.post( + "/api/telegram/webhook/secret-path", + json={"update_id": 123, "message": {"text": "/start"}}, + ) + + assert response.status_code == 200 + assert response.json() == {"ok": True} + assert calls == [[{"update_id": 123}]] + + +def test_telegram_webhook_rejects_wrong_secret(monkeypatch): + monkeypatch.setenv("TELEGRAM_WEBHOOK_SECRET", "secret-path") + + response = client.post( + "/api/telegram/webhook/wrong", + json={"update_id": 123}, + ) + + assert response.status_code == 403 + + +def test_start_bot_can_run_loops_without_polling(monkeypatch): + import src.bot.orchestrator as orchestrator + import src.database.db_manager as db_manager + import src.utils.config_loader as config_loader + + fake_bot = SimpleNamespace(poll_calls=0) + + def infinity_polling(**_kwargs): + fake_bot.poll_calls += 1 + + fake_bot.infinity_polling = infinity_polling + waited = [] + + class FakeCoordinator: + def __init__(self, **_kwargs): + pass + + def get_runtime_status(self): + return SimpleNamespace(loops=[]) + + def start_all(self): + return SimpleNamespace(loops=[]) + + monkeypatch.setenv("TELEGRAM_BOT_TOKEN", "123:test") + monkeypatch.setenv("TELEGRAM_BOT_POLLING_ENABLED", "false") + monkeypatch.setitem( + sys.modules, + "telebot", + SimpleNamespace(TeleBot=lambda _token: fake_bot), + ) + monkeypatch.setattr(orchestrator, "validate_or_raise", lambda _role: None) + monkeypatch.setattr(config_loader, "load_config", lambda: {}) + monkeypatch.setattr(db_manager, "DBManager", lambda: object()) + monkeypatch.setattr(orchestrator, "StartupCoordinator", FakeCoordinator) + monkeypatch.setattr(orchestrator, "_register_handlers", lambda **_kwargs: None) + monkeypatch.setattr(orchestrator, "_wait_forever", lambda: waited.append(True), raising=False) + + orchestrator.start_bot() + + assert fake_bot.poll_calls == 0 + assert waited == [True] diff --git a/web/app_factory.py b/web/app_factory.py index 5ea0066d..180d6488 100644 --- a/web/app_factory.py +++ b/web/app_factory.py @@ -20,6 +20,7 @@ from web.routers.payments import router as payments_router from web.routers.scan import router as scan_router from web.routers.sse_router import router as sse_router from web.routers.system import router as system_router +from web.routers.telegram_webhook import router as telegram_webhook_router from web.routes import router as legacy_router from web.scan_terminal_service import start_scan_terminal_prewarm @@ -58,6 +59,7 @@ def create_app() -> FastAPI: core_app.include_router(sse_router) core_app.include_router(payments_router) core_app.include_router(ops_router) + core_app.include_router(telegram_webhook_router) core_app.include_router(legacy_router) setattr(core_app.state, _ROUTES_REGISTERED_FLAG, True) if _scan_terminal_prewarm_enabled(): diff --git a/web/routers/telegram_webhook.py b/web/routers/telegram_webhook.py new file mode 100644 index 00000000..a796695d --- /dev/null +++ b/web/routers/telegram_webhook.py @@ -0,0 +1,12 @@ +"""Telegram webhook routes.""" + +from fastapi import APIRouter, Request + +from web.services.telegram_webhook import handle_telegram_webhook + +router = APIRouter(tags=["telegram"]) + + +@router.post("/api/telegram/webhook/{path_secret}") +async def telegram_webhook(path_secret: str, request: Request): + return await handle_telegram_webhook(path_secret, request) diff --git a/web/services/telegram_webhook.py b/web/services/telegram_webhook.py new file mode 100644 index 00000000..3d46bb04 --- /dev/null +++ b/web/services/telegram_webhook.py @@ -0,0 +1,108 @@ +from __future__ import annotations + +import hmac +import json +import os +import threading +from typing import Any + +from fastapi import HTTPException, Request +from loguru import logger + +_WEBHOOK_BOT: Any | None = None +_WEBHOOK_BOT_LOCK = threading.Lock() + + +def _configured_secret() -> str: + return str(os.getenv("TELEGRAM_WEBHOOK_SECRET") or "").strip() + + +def _configured_header_secret() -> str: + return str(os.getenv("TELEGRAM_WEBHOOK_HEADER_SECRET") or "").strip() + + +def _parse_update(payload: dict[str, Any]) -> Any: + from telebot import types # type: ignore + + return types.Update.de_json(json.dumps(payload)) + + +def _build_webhook_bot() -> Any: + import telebot # type: ignore + + from src.bot.handlers.activity import ActivityHandler + from src.bot.handlers.basic import BasicCommandHandler + from src.bot.io_layer import BotIOLayer + from src.bot.runtime_coordinator import StartupCoordinator + from src.database.db_manager import DBManager + from src.utils.config_loader import load_config + from src.utils.telegram_chat_ids import get_telegram_chat_ids_from_env + + token = str(os.getenv("TELEGRAM_BOT_TOKEN") or "").strip() + if not token: + raise RuntimeError("TELEGRAM_BOT_TOKEN is not configured") + + config = load_config() + bot = telebot.TeleBot(token) + io_layer = BotIOLayer(bot=bot, db=DBManager()) + startup_coordinator = StartupCoordinator( + bot=bot, + config=config, + command_access_mode="public", + protected_commands=[], + required_group_chat_id=",".join(get_telegram_chat_ids_from_env()), + ) + BasicCommandHandler( + bot=bot, + io_layer=io_layer, + runtime_status_provider=startup_coordinator.get_runtime_status, + config=config, + ).register() + ActivityHandler(bot=bot, io_layer=io_layer).register() + return bot + + +def _get_webhook_bot() -> Any: + global _WEBHOOK_BOT + if _WEBHOOK_BOT is not None: + return _WEBHOOK_BOT + with _WEBHOOK_BOT_LOCK: + if _WEBHOOK_BOT is None: + _WEBHOOK_BOT = _build_webhook_bot() + return _WEBHOOK_BOT + + +def _verify_webhook_secret(path_secret: str, request: Request) -> None: + expected = _configured_secret() + if not expected: + raise HTTPException(status_code=503, detail="telegram webhook is not configured") + if not hmac.compare_digest(str(path_secret or ""), expected): + raise HTTPException(status_code=403, detail="invalid telegram webhook secret") + + expected_header = _configured_header_secret() + if expected_header: + actual_header = request.headers.get("x-telegram-bot-api-secret-token", "") + if not hmac.compare_digest(actual_header, expected_header): + raise HTTPException(status_code=403, detail="invalid telegram webhook header") + + +async def handle_telegram_webhook(path_secret: str, request: Request) -> dict[str, bool]: + _verify_webhook_secret(path_secret, request) + try: + payload = await request.json() + except Exception as exc: + raise HTTPException(status_code=400, detail="invalid telegram webhook payload") from exc + if not isinstance(payload, dict): + raise HTTPException(status_code=400, detail="invalid telegram webhook payload") + + update = _parse_update(payload) + if update is None: + logger.warning("telegram webhook ignored unparsable update payload") + return {"ok": True} + + try: + _get_webhook_bot().process_new_updates([update]) + except Exception: + logger.exception("telegram webhook dispatch failed") + raise + return {"ok": True}