| import logging |
| import os |
| import hmac |
| from typing import Any, Dict, Optional |
| from fastapi import Request |
| import httpx |
|
|
| from core.communication.adapters.base import PlatformAdapter |
|
|
| logger = logging.getLogger(__name__) |
|
|
| class TelegramAdapter(PlatformAdapter): |
| """ |
| Adapter for Telegram Bot API. |
| """ |
| def __init__( |
| self, |
| bot_token: Optional[str] = None, |
| secret_token: Optional[str] = None, |
| ): |
| self.bot_token = bot_token or os.getenv("TELEGRAM_BOT_TOKEN") |
| self.secret_token = secret_token or os.getenv("TELEGRAM_SECRET_TOKEN") |
| self.api_base = f"https://api.telegram.org/bot{self.bot_token}" |
|
|
| async def verify_request(self, request: Request, body_bytes: bytes) -> bool: |
| """ |
| Verify Telegram webhook request using constant-time comparison. |
| |
| Uses hmac.compare_digest() to prevent timing attacks where attackers |
| measure response times to guess valid tokens character-by-character. |
| """ |
| if not self.secret_token: |
| |
| |
| return True |
|
|
| header_token = request.headers.get("X-Telegram-Bot-Api-Secret-Token") |
| return hmac.compare_digest(header_token, self.secret_token) |
|
|
| async def normalize_payload(self, request: Request, body_bytes: bytes) -> Dict[str, Any]: |
| """ |
| Normalize standard Telegram JSON update. |
| """ |
| import json |
| try: |
| data = json.loads(body_bytes) |
| except json.JSONDecodeError: |
| return {} |
| |
| |
| message = data.get("message", {}) |
| |
| user_id = str(message.get("from", {}).get("id", "")) |
| chat_id = str(message.get("chat", {}).get("id", "")) |
| text = message.get("text", "") |
| |
| |
| if "voice" in message: |
| text = "[Voice Message]" |
| data["media_id"] = message["voice"]["file_id"] |
| data["media_type"] = "voice" |
| |
| |
| username = message.get("from", {}).get("username", "") |
| |
| return { |
| "platform": "telegram", |
| "user_id": user_id, |
| "username": username, |
| "channel_id": chat_id, |
| "content": text, |
| "metadata": data |
| } |
|
|
| async def send_message(self, target_id: str, message: str, metadata: Dict = None) -> bool: |
| if not self.bot_token: |
| logger.error("Telegram token not set.") |
| return False |
| |
| try: |
| async with httpx.AsyncClient() as client: |
| payload = { |
| "chat_id": target_id, |
| "text": message, |
| "parse_mode": "Markdown" |
| } |
| response = await client.post(f"{self.api_base}/sendMessage", json=payload) |
| response.raise_for_status() |
| return True |
| except Exception as e: |
| logger.error(f"Failed to send Telegram message: {e}") |
| return False |
|
|
| async def get_media(self, media_id: str) -> Optional[bytes]: |
| """Downloads voice/file from Telegram.""" |
| if not self.bot_token: |
| return None |
| try: |
| async with httpx.AsyncClient() as client: |
| |
| res = await client.get(f"{self.api_base}/getFile", params={"file_id": media_id}) |
| res.raise_for_status() |
| file_path = res.json().get("result", {}).get("file_path") |
| if not file_path: |
| return None |
| |
| |
| download_url = f"https://api.telegram.org/file/bot{self.bot_token}/{file_path}" |
| media_res = await client.get(download_url) |
| media_res.raise_for_status() |
| return media_res.content |
| except Exception as e: |
| logger.error(f"Telegram media download failed: {e}") |
| return None |
|
|