Add Telegram webhook command entrypoint
This commit is contained in:
@@ -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
|
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(
|
def _register_handlers(
|
||||||
bot: Any,
|
bot: Any,
|
||||||
config: dict[str, Any],
|
config: dict[str, Any],
|
||||||
@@ -73,4 +87,8 @@ def start_bot() -> None:
|
|||||||
started_count,
|
started_count,
|
||||||
len(runtime_status.loops),
|
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"])
|
bot.infinity_polling(allowed_updates=["message", "callback_query", "chat_join_request"])
|
||||||
|
|||||||
@@ -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]
|
||||||
@@ -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.scan import router as scan_router
|
||||||
from web.routers.sse_router import router as sse_router
|
from web.routers.sse_router import router as sse_router
|
||||||
from web.routers.system import router as system_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.routes import router as legacy_router
|
||||||
from web.scan_terminal_service import start_scan_terminal_prewarm
|
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(sse_router)
|
||||||
core_app.include_router(payments_router)
|
core_app.include_router(payments_router)
|
||||||
core_app.include_router(ops_router)
|
core_app.include_router(ops_router)
|
||||||
|
core_app.include_router(telegram_webhook_router)
|
||||||
core_app.include_router(legacy_router)
|
core_app.include_router(legacy_router)
|
||||||
setattr(core_app.state, _ROUTES_REGISTERED_FLAG, True)
|
setattr(core_app.state, _ROUTES_REGISTERED_FLAG, True)
|
||||||
if _scan_terminal_prewarm_enabled():
|
if _scan_terminal_prewarm_enabled():
|
||||||
|
|||||||
@@ -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)
|
||||||
@@ -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}
|
||||||
Reference in New Issue
Block a user