This commit is contained in:
kislovdm
2026-04-08 23:56:38 +03:00
commit 3ca8f7545f
12 changed files with 620 additions and 0 deletions
+172
View File
@@ -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))
+28
View File
@@ -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",
)
+56
View File
@@ -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()
+110
View File
@@ -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,
)
+12
View File
@@ -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)
+3
View File
@@ -0,0 +1,3 @@
git+https://github.com/MaxApiTeam/PyMax.git@dev/1.2.6
requests
python-dotenv
+41
View File
@@ -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()
+130
View File
@@ -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