commit 3ca8f7545fc1f114f448eb9a75e74a0705da5efe Author: kislovdm Date: Wed Apr 8 23:56:38 2026 +0300 init diff --git a/.env.sample b/.env.sample new file mode 100644 index 0000000..7853129 --- /dev/null +++ b/.env.sample @@ -0,0 +1,7 @@ +TZ=Europe/Moscow +PYTHONUNBUFFERED=0 +MAX_PHONE=+10000000000 +MAX_WORK_DIR=cache +TELEGRAM_BOT_TOKEN=123456789:your_bot_token +TELEGRAM_FALLBACK_USER_ID=123456789 +SQLITE_PATH=max2telegram.db diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..192af1f --- /dev/null +++ b/.gitignore @@ -0,0 +1,34 @@ +.env +cache/* + +# Python bytecode / cache +__pycache__/ +*.py[cod] +*$py.class + +# Virtual environments +.venv/ +venv/ +env/ +ENV/ + +# Build / packaging +build/ +dist/ +*.egg-info/ +.eggs/ +pip-wheel-metadata/ + +# Tool caches +.pytest_cache/ +.mypy_cache/ +.ruff_cache/ +.coverage +.coverage.* +htmlcov/ + +# IDE / OS +.idea/ +.vscode/ +.DS_Store +Thumbs.db \ No newline at end of file diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..94859ab --- /dev/null +++ b/Dockerfile @@ -0,0 +1,10 @@ +FROM python:3.12-slim + +WORKDIR /app + +RUN apt-get update; apt-get install -y git +COPY ./src/requirements.txt . +RUN pip install --no-cache-dir -r ./requirements.txt + +COPY ./src . +CMD [ "python", "-u", "./main.py" ] diff --git a/docker-compose.yaml b/docker-compose.yaml new file mode 100644 index 0000000..426645c --- /dev/null +++ b/docker-compose.yaml @@ -0,0 +1,17 @@ +version: '3' +services: + max: + build: ./ + container_name: max + restart: always + env_file: + - ./.env + volumes: + - "./cache:/app/cache" + - "/etc/hosts:/etc/hosts" + ports: + - 5004:5000 + logging: + driver: json-file + options: + max-size: 50m diff --git a/src/bridge.py b/src/bridge.py new file mode 100644 index 0000000..d2ca1b6 --- /dev/null +++ b/src/bridge.py @@ -0,0 +1,172 @@ +import logging +from typing import Any + +from max_parser import parse_message +from models import ParsedMessage +from pymax import MaxClient +from pymax.types import PhotoAttach, VideoAttach +from storage import BridgeStorage +from telegram_api import TelegramClient + +logger = logging.getLogger(__name__) + + +class MaxToTelegramBridge: + def __init__(self, max_client: MaxClient, telegram: TelegramClient, storage: BridgeStorage) -> None: + self._max_client = max_client + self._telegram = telegram + self._storage = storage + + async def forward_message(self, max_message: Any) -> None: + parsed = parse_message(max_message) + parsed = await self._enrich_from_max(max_message, parsed) + + if self._storage.was_forwarded(parsed.message_id, parsed.chat_id): + logger.debug("Skip duplicated message %s/%s", parsed.chat_id, parsed.message_id) + return + + text = self._format_caption(parsed.sender_name, parsed.chat_name, parsed.text) + has_any_payload = bool(text.strip()) or bool(parsed.image_urls) or bool(parsed.video_urls) + if not has_any_payload: + logger.debug("Skip empty message %s/%s", parsed.chat_id, parsed.message_id) + return + + target_chat_id, matched_by_title = await self._telegram.resolve_target_chat_id(parsed.chat_name) + if matched_by_title: + logger.info("Route Max chat '%s' to Telegram chat %s", parsed.chat_name, target_chat_id) + else: + logger.info( + "Telegram chat '%s' not found, route to fallback user %s", + parsed.chat_name, + target_chat_id, + ) + + total_media = len(parsed.image_urls) + len(parsed.video_urls) + if total_media > 1: + # Отправляем одним альбомом в Telegram (единое сообщение). + await self._telegram.send_media_group( + chat_id=target_chat_id, + image_urls=parsed.image_urls, + video_urls=parsed.video_urls, + caption=text, + ) + self._storage.mark_forwarded(parsed.message_id, parsed.chat_id) + logger.info( + "Forwarded media group %s/%s (images=%s, videos=%s)", + parsed.chat_id, + parsed.message_id, + len(parsed.image_urls), + len(parsed.video_urls), + ) + return + + sent_any = False + if parsed.text.strip() and total_media == 0: + await self._telegram.send_text(target_chat_id, text) + sent_any = True + + for index, image_url in enumerate(parsed.image_urls): + caption = text if not sent_any and index == 0 else None + await self._telegram.send_photo(target_chat_id, image_url, caption=caption) + sent_any = True + + for index, video_url in enumerate(parsed.video_urls): + caption = text if not sent_any and index == 0 else None + await self._telegram.send_video(target_chat_id, video_url, caption=caption) + sent_any = True + + self._storage.mark_forwarded(parsed.message_id, parsed.chat_id) + logger.info( + "Forwarded message %s/%s (images=%s, videos=%s)", + parsed.chat_id, + parsed.message_id, + len(parsed.image_urls), + len(parsed.video_urls), + ) + + @staticmethod + def _format_caption(sender_name: str, chat_name: str, text: str) -> str: + header = f"MAX: {sender_name} / {chat_name}" + if text.strip(): + return f"{header}\n\n{text}" + return header + + async def _enrich_from_max(self, max_message: Any, parsed: ParsedMessage) -> ParsedMessage: + # Имена отправителя и чата берем из API Max, чтобы всегда получить человекочитаемый формат. + try: + user = await self._max_client.get_user(user_id=max_message.sender) + if user and getattr(user, "names", None): + first_name = getattr(user.names[0], "name", "") + if first_name: + parsed.sender_name = str(first_name) + except Exception: + logger.debug("Cannot resolve sender name", exc_info=True) + + try: + chat = await self._max_client.get_chat(chat_id=max_message.chat_id) + title = getattr(chat, "title", None) + if title: + parsed.chat_name = str(title) + except Exception: + logger.debug("Cannot resolve chat title", exc_info=True) + + attaches = getattr(max_message, "attaches", None) or [] + for attach in attaches: + if isinstance(attach, PhotoAttach): + parsed.image_urls.extend(self._extract_photo_urls(attach)) + elif isinstance(attach, VideoAttach): + try: + video = await self._max_client.get_video_by_id( + chat_id=max_message.chat_id, + message_id=max_message.id, + video_id=attach.video_id, + ) + video_url = getattr(video, "url", None) + if video_url: + parsed.video_urls.append(str(video_url)) + except Exception: + logger.exception("Cannot resolve video URL from Max") + + # Убираем дубли URL, если парсер и enrich нашли одинаковые вложения. + parsed.image_urls = list(dict.fromkeys(parsed.image_urls)) + parsed.video_urls = list(dict.fromkeys(parsed.video_urls)) + return parsed + + def _extract_photo_urls(self, attach: PhotoAttach) -> list[str]: + urls: list[str] = [] + seen_ids: set[int] = set() + + def walk(node: Any) -> None: + if node is None: + return + + obj_id = id(node) + if obj_id in seen_ids: + return + seen_ids.add(obj_id) + + if isinstance(node, str): + if node.startswith("http://") or node.startswith("https://"): + urls.append(node) + return + + if isinstance(node, (list, tuple, set)): + for item in node: + walk(item) + return + + if isinstance(node, dict): + for key, value in node.items(): + # Поиск всех возможных URL полей, включая альбомы/варианты размеров. + if key in {"base_url", "url", "src", "download_url"} and isinstance(value, str): + if value.startswith("http://") or value.startswith("https://"): + urls.append(value) + else: + walk(value) + return + + if hasattr(node, "__dict__"): + walk(vars(node)) + + walk(attach) + return list(dict.fromkeys(urls)) diff --git a/src/config.py b/src/config.py new file mode 100644 index 0000000..9cda0bd --- /dev/null +++ b/src/config.py @@ -0,0 +1,28 @@ +import os +from dataclasses import dataclass + + +@dataclass(frozen=True) +class Settings: + max_phone: str + max_work_dir: str + telegram_bot_token: str + telegram_fallback_user_id: str + sqlite_path: str + + +def _require_env(name: str) -> str: + value = os.getenv(name, "").strip() + if not value: + raise ValueError(f"Environment variable {name} is required") + return value + + +def load_settings() -> Settings: + return Settings( + max_phone=_require_env("MAX_PHONE"), + max_work_dir=os.getenv("MAX_WORK_DIR", "cache").strip() or "cache", + telegram_bot_token=_require_env("TELEGRAM_BOT_TOKEN"), + telegram_fallback_user_id=_require_env("TELEGRAM_FALLBACK_USER_ID"), + sqlite_path=os.getenv("SQLITE_PATH", "max2telegram.db").strip() or "max2telegram.db", + ) diff --git a/src/main.py b/src/main.py new file mode 100644 index 0000000..7610131 --- /dev/null +++ b/src/main.py @@ -0,0 +1,56 @@ +import asyncio +import logging + +from pymax import MaxClient, Message +from dotenv import load_dotenv + +from bridge import MaxToTelegramBridge +from config import load_settings +from storage import BridgeStorage +from telegram_api import TelegramClient + + +def _setup_logging() -> None: + logging.basicConfig( + level=logging.INFO, + format="%(asctime)s %(levelname)s %(name)s: %(message)s", + ) + + +def build_client() -> tuple[MaxClient, MaxToTelegramBridge]: + settings = load_settings() + max_client = MaxClient( + phone=settings.max_phone, + work_dir=settings.max_work_dir, + ) + telegram_client = TelegramClient( + bot_token=settings.telegram_bot_token, + fallback_user_id=settings.telegram_fallback_user_id, + ) + storage = BridgeStorage(settings.sqlite_path) + bridge = MaxToTelegramBridge(max_client=max_client, telegram=telegram_client, storage=storage) + return max_client, bridge + + +def main() -> None: + load_dotenv() + _setup_logging() + max_client, bridge = build_client() + logger = logging.getLogger("max2telegram") + + @max_client.on_start + async def on_start() -> None: + logger.info("Max client started as %s", max_client.me.id) + + @max_client.on_message() + async def on_message(message: Message) -> None: + try: + await bridge.forward_message(message) + except Exception: + logger.exception("Failed to forward Max message") + + asyncio.run(max_client.start()) + + +if __name__ == "__main__": + main() diff --git a/src/max_parser.py b/src/max_parser.py new file mode 100644 index 0000000..fb4b220 --- /dev/null +++ b/src/max_parser.py @@ -0,0 +1,110 @@ +from typing import Any + +from models import ParsedMessage + + +def _stringify(value: Any) -> str: + if value is None: + return "" + return str(value).strip() + + +def _get_attr(obj: Any, names: list[str], default: Any = None) -> Any: + for name in names: + if hasattr(obj, name): + value = getattr(obj, name) + if value is not None: + return value + return default + + +def _as_dict(obj: Any) -> dict[str, Any]: + if isinstance(obj, dict): + return obj + if hasattr(obj, "__dict__"): + return vars(obj) + return {} + + +def _is_image(media_type: str) -> bool: + value = media_type.lower() + return "image" in value or "photo" in value or value in {"jpg", "jpeg", "png", "webp"} + + +def _is_video(media_type: str) -> bool: + value = media_type.lower() + return "video" in value or value in {"mp4", "mov", "mkv", "avi"} + + +def _extract_media_urls(message: Any) -> tuple[list[str], list[str]]: + image_urls: list[str] = [] + video_urls: list[str] = [] + + # В PyMax рабочее поле для вложений обычно называется attaches. + raw_attachments = _get_attr(message, ["attaches", "attachments", "media", "files"], default=[]) or [] + for item in raw_attachments: + data = _as_dict(item) + media_type = _stringify(data.get("type") or data.get("media_type") or data.get("kind")) + url = _stringify( + data.get("base_url") + or data.get("url") + or data.get("link") + or data.get("download_url") + or data.get("src") + ) + + if not url: + nested = data.get("file") or data.get("payload") + nested_data = _as_dict(nested) + url = _stringify( + nested_data.get("base_url") + or nested_data.get("url") + or nested_data.get("link") + or nested_data.get("download_url") + or nested_data.get("src") + ) + if not media_type: + media_type = _stringify(nested_data.get("type") or nested_data.get("media_type")) + + if not url: + continue + + if _is_image(media_type): + image_urls.append(url) + elif _is_video(media_type): + video_urls.append(url) + + return image_urls, video_urls + + +def parse_message(message: Any) -> ParsedMessage: + sender = _get_attr(message, ["sender", "sender_name", "author"], default="unknown") + sender_data = _as_dict(sender) + sender_name = ( + _stringify(sender_data.get("nickname")) + or _stringify(sender_data.get("username")) + or _stringify(sender_data.get("name")) + or _stringify(sender) + or "unknown" + ) + + chat_name = ( + _stringify(_get_attr(message, ["chat_title", "chat_name", "group_name"])) + or _stringify(_get_attr(message, ["chat"], default="")) + or "direct" + ) + + message_id = _stringify(_get_attr(message, ["id", "message_id", "mid"])) or "unknown-id" + chat_id = _stringify(_get_attr(message, ["chat_id", "dialog_id", "peer_id"])) or "unknown-chat" + text = _stringify(_get_attr(message, ["text", "message", "body"])) + + image_urls, video_urls = _extract_media_urls(message) + return ParsedMessage( + message_id=message_id, + chat_id=chat_id, + sender_name=sender_name, + chat_name=chat_name, + text=text, + image_urls=image_urls, + video_urls=video_urls, + ) diff --git a/src/models.py b/src/models.py new file mode 100644 index 0000000..53a8be5 --- /dev/null +++ b/src/models.py @@ -0,0 +1,12 @@ +from dataclasses import dataclass, field + + +@dataclass +class ParsedMessage: + message_id: str + chat_id: str + sender_name: str + chat_name: str + text: str + image_urls: list[str] = field(default_factory=list) + video_urls: list[str] = field(default_factory=list) diff --git a/src/requirements.txt b/src/requirements.txt new file mode 100644 index 0000000..a16ef6a --- /dev/null +++ b/src/requirements.txt @@ -0,0 +1,3 @@ +git+https://github.com/MaxApiTeam/PyMax.git@dev/1.2.6 +requests +python-dotenv diff --git a/src/storage.py b/src/storage.py new file mode 100644 index 0000000..34a65c5 --- /dev/null +++ b/src/storage.py @@ -0,0 +1,41 @@ +import sqlite3 +from contextlib import closing + + +class BridgeStorage: + def __init__(self, db_path: str) -> None: + self._db_path = db_path + self._init_db() + + def _connect(self) -> sqlite3.Connection: + return sqlite3.connect(self._db_path) + + def _init_db(self) -> None: + with closing(self._connect()) as conn: + conn.execute( + """ + CREATE TABLE IF NOT EXISTS forwarded_messages ( + message_id TEXT NOT NULL, + chat_id TEXT NOT NULL, + forwarded_at DATETIME DEFAULT CURRENT_TIMESTAMP, + PRIMARY KEY (message_id, chat_id) + ) + """ + ) + conn.commit() + + def was_forwarded(self, message_id: str, chat_id: str) -> bool: + with closing(self._connect()) as conn: + row = conn.execute( + "SELECT 1 FROM forwarded_messages WHERE message_id = ? AND chat_id = ?", + (message_id, chat_id), + ).fetchone() + return row is not None + + def mark_forwarded(self, message_id: str, chat_id: str) -> None: + with closing(self._connect()) as conn: + conn.execute( + "INSERT OR IGNORE INTO forwarded_messages (message_id, chat_id) VALUES (?, ?)", + (message_id, chat_id), + ) + conn.commit() diff --git a/src/telegram_api.py b/src/telegram_api.py new file mode 100644 index 0000000..f5a17d6 --- /dev/null +++ b/src/telegram_api.py @@ -0,0 +1,130 @@ +import asyncio +from typing import Any + +import requests + + +class TelegramApiError(RuntimeError): + pass + + +class TelegramClient: + def __init__(self, bot_token: str, fallback_user_id: str, timeout: int = 30) -> None: + self._base_url = f"https://api.telegram.org/bot{bot_token}" + self._fallback_user_id = fallback_user_id + self._timeout = timeout + self._chat_title_to_id: dict[str, str] = {} + + async def resolve_target_chat_id(self, max_chat_name: str) -> tuple[str, bool]: + chat_id = await self._find_chat_id_by_title(max_chat_name) + if chat_id: + return chat_id, True + return self._fallback_user_id, False + + async def send_text(self, chat_id: str, text: str) -> None: + await self._request( + "sendMessage", + { + "chat_id": chat_id, + "text": text, + "disable_web_page_preview": True, + }, + ) + + async def send_photo(self, chat_id: str, photo_url: str, caption: str | None = None) -> None: + payload: dict[str, Any] = { + "chat_id": chat_id, + "photo": photo_url, + } + if caption: + payload["caption"] = caption + await self._request("sendPhoto", payload) + + async def send_video(self, chat_id: str, video_url: str, caption: str | None = None) -> None: + payload: dict[str, Any] = { + "chat_id": chat_id, + "video": video_url, + "supports_streaming": True, + } + if caption: + payload["caption"] = caption + await self._request("sendVideo", payload) + + async def send_media_group( + self, + chat_id: str, + image_urls: list[str], + video_urls: list[str], + caption: str | None = None, + ) -> None: + media: list[dict[str, Any]] = [] + + for url in image_urls: + media.append({"type": "photo", "media": url}) + + for url in video_urls: + media.append({"type": "video", "media": url, "supports_streaming": True}) + + if not media: + return + + if caption: + media[0]["caption"] = caption + + await self._request( + "sendMediaGroup", + { + "chat_id": chat_id, + "media": media, + }, + ) + + async def _find_chat_id_by_title(self, chat_title: str) -> str | None: + normalized = self._normalize_title(chat_title) + if not normalized: + return None + + cached = self._chat_title_to_id.get(normalized) + if cached: + return cached + + response = await self._request("getUpdates", {"timeout": 0, "limit": 100}) + for update in response.get("result", []): + for container in ("message", "edited_message", "channel_post", "edited_channel_post"): + message = update.get(container) + if not isinstance(message, dict): + continue + chat = message.get("chat") + if not isinstance(chat, dict): + continue + + title_value = self._extract_chat_title(chat) + chat_id = chat.get("id") + if title_value and chat_id is not None: + self._chat_title_to_id[self._normalize_title(title_value)] = str(chat_id) + + return self._chat_title_to_id.get(normalized) + + @staticmethod + def _extract_chat_title(chat: dict[str, Any]) -> str: + return str(chat.get("title") or chat.get("username") or "").strip() + + @staticmethod + def _normalize_title(value: str) -> str: + return value.strip().casefold() + + async def _request(self, method: str, payload: dict[str, Any]) -> dict[str, Any]: + url = f"{self._base_url}/{method}" + + def _do_request() -> requests.Response: + return requests.post(url, json=payload, timeout=self._timeout) + + response = await asyncio.to_thread(_do_request) + if response.status_code >= 400: + raise TelegramApiError( + f"Telegram HTTP error on {method}: {response.status_code} {response.text}" + ) + data = response.json() + if not data.get("ok"): + raise TelegramApiError(f"Telegram API error on {method}: {data}") + return data