7 Commits
Author SHA1 Message Date
kislovdm d5ddde3103 ++
Docker Hub / build-and-push (push) Failing after 13s
2026-04-22 14:22:59 +03:00
kislovdm cabfb9b5f9 ++ 2026-04-22 13:25:49 +03:00
kislovdm cc9d05b11b ++ 2026-04-22 12:54:14 +03:00
kislovdm 64040d6bc5 ++ 2026-04-22 12:44:22 +03:00
kislovdm 87c4c765e5 tech specs 2026-04-21 00:10:19 +03:00
kislovdm 2d32bf4807 ++
Docker Hub / build-and-push (push) Failing after 13s
2026-04-20 19:19:45 +03:00
kislovdm d019e884b6 обработка forward 2026-04-20 19:13:08 +03:00
5 changed files with 839 additions and 24 deletions
+180 -15
View File
@@ -5,7 +5,7 @@ from typing import Any
from max_parser import parse_message from max_parser import parse_message
from models import ParsedMessage from models import ParsedMessage
from pymax import MaxClient from pymax import MaxClient
from pymax.types import PhotoAttach, VideoAttach from pymax.types import AudioAttach, FileAttach, Message, PhotoAttach, StickerAttach, VideoAttach
from storage import BridgeStorage from storage import BridgeStorage
from telegram_api import TelegramApiError, TelegramClient from telegram_api import TelegramApiError, TelegramClient
@@ -110,7 +110,8 @@ class MaxToTelegramBridge:
) )
return return
if parsed.text.strip() and total_media == 0: should_send_plain_text = total_media == 0 and not parsed.file_urls and bool(text.strip())
if should_send_plain_text:
target_chat_id, sent = await self._send_with_migration_retry( target_chat_id, sent = await self._send_with_migration_retry(
target_chat_id=target_chat_id, target_chat_id=target_chat_id,
max_chat_title_norm=normalized, max_chat_title_norm=normalized,
@@ -177,6 +178,7 @@ class MaxToTelegramBridge:
for index, file_url in enumerate(parsed.file_urls): for index, file_url in enumerate(parsed.file_urls):
caption = text if not sent_any and index == 0 else None caption = text if not sent_any and index == 0 else None
file_name = parsed.file_names_by_url.get(file_url)
target_chat_id, sent = await self._send_with_migration_retry( target_chat_id, sent = await self._send_with_migration_retry(
target_chat_id=target_chat_id, target_chat_id=target_chat_id,
max_chat_title_norm=normalized, max_chat_title_norm=normalized,
@@ -184,6 +186,7 @@ class MaxToTelegramBridge:
send_action=lambda chat_id: self._telegram.send_document( send_action=lambda chat_id: self._telegram.send_document(
chat_id, chat_id,
file_url, file_url,
file_name=file_name,
caption=caption, caption=caption,
reply_to_message_id=reply_telegram_mid if not sent_any and index == 0 else None, reply_to_message_id=reply_telegram_mid if not sent_any and index == 0 else None,
), ),
@@ -200,12 +203,13 @@ class MaxToTelegramBridge:
if not sent_any: if not sent_any:
# Последняя страховка: гарантируем уведомление в Telegram даже для пустых/неизвестных payload. # Последняя страховка: гарантируем уведомление в Telegram даже для пустых/неизвестных payload.
fallback_text = text.strip() or self._build_fallback_unknown_notice(parsed)
target_chat_id, sent = await self._send_with_migration_retry( target_chat_id, sent = await self._send_with_migration_retry(
target_chat_id=target_chat_id, target_chat_id=target_chat_id,
max_chat_title_norm=normalized, max_chat_title_norm=normalized,
max_chat_title=parsed.chat_name, max_chat_title=parsed.chat_name,
send_action=lambda chat_id: self._telegram.send_text( send_action=lambda chat_id: self._telegram.send_text(
chat_id, self._build_fallback_unknown_notice(parsed), reply_to_message_id=reply_telegram_mid chat_id, fallback_text, reply_to_message_id=reply_telegram_mid
), ),
) )
mid = sent.get("result", {}).get("message_id") if isinstance(sent.get("result"), dict) else None mid = sent.get("result", {}).get("message_id") if isinstance(sent.get("result"), dict) else None
@@ -338,35 +342,172 @@ class MaxToTelegramBridge:
except Exception: except Exception:
logger.debug("Cannot resolve chat title", exc_info=True) logger.debug("Cannot resolve chat title", exc_info=True)
attaches = getattr(max_message, "attaches", None) or [] await self._collect_message_attachments(
message=max_message,
parsed=parsed,
source_tag="root",
)
link = getattr(max_message, "link", None)
linked_message = getattr(link, "message", None)
if linked_message is not None:
# Для reply не копируем текст исходного сообщения в тело:
# иначе получаем дубль (цитата + тот же текст как новое сообщение).
is_reply = bool(parsed.reply_to_max_message_id)
if not is_reply and not (parsed.text or "").strip():
linked_text = str(getattr(linked_message, "text", "") or "").strip()
if linked_text:
parsed.text = linked_text
# Для reply нельзя переносить вложения linked_message:
# это исходное сообщение, и его медиа не должны отправляться повторно.
if not is_reply:
await self._collect_message_attachments(
message=linked_message,
parsed=parsed,
source_tag="forward",
)
# Убираем дубли URL, если парсер и enrich нашли одинаковые вложения.
parsed.image_urls = list(dict.fromkeys(parsed.image_urls))
parsed.video_urls = list(dict.fromkeys(parsed.video_urls))
parsed.file_urls = list(dict.fromkeys(parsed.file_urls))
parsed.file_names_by_url = {
url: name for url, name in parsed.file_names_by_url.items() if url in parsed.file_urls and name
}
if parsed.file_urls:
parsed.unknown_attachments = [
x for x in parsed.unknown_attachments if not self._is_file_unknown_marker(x)
]
parsed.unknown_attachments = list(dict.fromkeys(parsed.unknown_attachments))
return parsed
async def _collect_message_attachments(
self,
*,
message: Any,
parsed: ParsedMessage,
source_tag: str,
) -> None:
attaches = getattr(message, "attaches", None) or []
message_chat_id = getattr(message, "chat_id", None)
message_id = getattr(message, "id", None)
if message_chat_id is None:
message_chat_id = parsed.chat_id
for attach in attaches: for attach in attaches:
if isinstance(attach, PhotoAttach): if isinstance(attach, PhotoAttach):
parsed.image_urls.extend(self._extract_photo_urls(attach)) parsed.image_urls.extend(self._extract_photo_urls(attach))
elif isinstance(attach, VideoAttach): continue
if isinstance(attach, VideoAttach):
try: try:
video = await self._max_client.get_video_by_id( video = await self._max_client.get_video_by_id(
chat_id=max_message.chat_id, chat_id=message_chat_id,
message_id=max_message.id, message_id=message_id,
video_id=attach.video_id, video_id=attach.video_id,
) )
video_url = getattr(video, "url", None) video_url = getattr(video, "url", None)
if video_url: if video_url:
parsed.video_urls.append(str(video_url)) parsed.video_urls.append(str(video_url))
except Exception: except Exception:
logger.exception("Cannot resolve video URL from Max") logger.exception("Cannot resolve video URL from Max (%s)", source_tag)
continue
if isinstance(attach, FileAttach):
resolved = await self._resolve_file_attach_url(
message_chat_id=message_chat_id,
message_id=message_id,
attach=attach,
)
if resolved:
parsed.file_urls.append(resolved)
file_name = str(getattr(attach, "name", "") or "").strip()
if file_name:
parsed.file_names_by_url[resolved] = file_name
else: else:
fallback = str(getattr(attach, "name", "") or "").strip()
if fallback:
parsed.text = self._append_missing_file_note(parsed.text, fallback)
else:
parsed.unknown_attachments.append(type(attach).__name__)
continue
if isinstance(attach, AudioAttach):
audio_url = str(getattr(attach, "url", "") or "").strip()
if audio_url:
parsed.file_urls.append(audio_url)
else:
parsed.unknown_attachments.append(type(attach).__name__)
continue
if isinstance(attach, StickerAttach):
sticker_url = str(getattr(attach, "url", "") or "").strip()
if sticker_url:
parsed.image_urls.append(sticker_url)
else:
parsed.unknown_attachments.append(type(attach).__name__)
continue
# Fallback на случай сырого Attach/нестандартного типа:
if await self._resolve_generic_file_attach(
message_chat_id=message_chat_id,
message_id=message_id,
attach=attach,
parsed=parsed,
):
continue
urls = self._extract_any_urls(attach) urls = self._extract_any_urls(attach)
if urls: if urls:
parsed.file_urls.extend(urls) parsed.file_urls.extend(urls)
else: continue
parsed.unknown_attachments.append(type(attach).__name__) parsed.unknown_attachments.append(type(attach).__name__)
# Убираем дубли URL, если парсер и enrich нашли одинаковые вложения. async def _resolve_file_attach_url(self, *, message_chat_id: Any, message_id: Any, attach: FileAttach) -> str | None:
parsed.image_urls = list(dict.fromkeys(parsed.image_urls)) file_id = getattr(attach, "file_id", None)
parsed.video_urls = list(dict.fromkeys(parsed.video_urls)) if file_id is None or message_id is None:
parsed.file_urls = list(dict.fromkeys(parsed.file_urls)) return None
parsed.unknown_attachments = list(dict.fromkeys(parsed.unknown_attachments)) try:
return parsed file_info = await self._max_client.get_file_by_id(
chat_id=message_chat_id,
message_id=message_id,
file_id=file_id,
)
url = getattr(file_info, "url", None)
return str(url) if url else None
except Exception:
logger.exception("Cannot resolve file URL from Max (file_id=%s)", file_id)
return None
async def _resolve_generic_file_attach(
self,
*,
message_chat_id: Any,
message_id: Any,
attach: Any,
parsed: ParsedMessage,
) -> bool:
file_id = getattr(attach, "file_id", None)
if file_id is None or message_id is None:
return False
try:
file_info = await self._max_client.get_file_by_id(
chat_id=message_chat_id,
message_id=message_id,
file_id=file_id,
)
url = getattr(file_info, "url", None)
if url:
resolved = str(url)
parsed.file_urls.append(resolved)
file_name = str(getattr(attach, "name", "") or "").strip()
if file_name:
parsed.file_names_by_url[resolved] = file_name
return True
except Exception:
logger.debug("Cannot resolve generic file attach from Max", exc_info=True)
return False
def _is_self_message(self, max_message: Any) -> bool: def _is_self_message(self, max_message: Any) -> bool:
sender = getattr(max_message, "sender", None) sender = getattr(max_message, "sender", None)
@@ -445,6 +586,25 @@ class MaxToTelegramBridge:
walk(node) walk(node)
return list(dict.fromkeys(urls)) return list(dict.fromkeys(urls))
@staticmethod
def _is_forward_attach_like(attach: Any) -> bool:
name = type(attach).__name__.lower()
if "forward" in name or "share" in name or "quote" in name:
return True
if hasattr(attach, "__dict__"):
keys = {str(k).lower() for k in vars(attach).keys()}
if {"forward", "forwarded", "forwards", "link", "message", "messages", "origin", "payload"} & keys:
return True
return False
@staticmethod
def _append_missing_file_note(current_text: str, file_name: str) -> str:
text = (current_text or "").strip()
note = f"[MAX forwarded file without direct URL] {file_name}"
if not text:
return note
return f"{text}\n{note}"
@staticmethod @staticmethod
def _append_unknown_attachment_notice(*, parsed: ParsedMessage, text: str) -> str: def _append_unknown_attachment_notice(*, parsed: ParsedMessage, text: str) -> str:
if not parsed.unknown_attachments: if not parsed.unknown_attachments:
@@ -463,3 +623,8 @@ class MaxToTelegramBridge:
) )
unknown = ", ".join(parsed.unknown_attachments[:5]) if parsed.unknown_attachments else "unknown" unknown = ", ".join(parsed.unknown_attachments[:5]) if parsed.unknown_attachments else "unknown"
return f"{base}\n\n[!] Неизвестный или пустой тип сообщения из MAX (attachments={unknown})." return f"{base}\n\n[!] Неизвестный или пустой тип сообщения из MAX (attachments={unknown})."
@staticmethod
def _is_file_unknown_marker(value: str) -> bool:
normalized = str(value or "").strip().casefold()
return normalized in {"attachtype.file", "fileattach", "file"}
+76
View File
@@ -36,6 +36,64 @@ def _is_video(media_type: str) -> bool:
return "video" in value or value in {"mp4", "mov", "mkv", "avi"} return "video" in value or value in {"mp4", "mov", "mkv", "avi"}
def _is_forward_like(data: dict[str, Any]) -> bool:
media_type = _stringify(data.get("type") or data.get("media_type") or data.get("kind")).lower()
if "forward" in media_type or "share" in media_type or "quote" in media_type:
return True
forward_keys = {
"forward",
"forwarded",
"forwards",
"link",
"message",
"messages",
"payload",
"quote",
"origin",
}
return any(key in data for key in forward_keys)
def _classify_url(url: str, media_type: str) -> str:
lowered = url.lower()
if _is_image(media_type) or lowered.endswith((".jpg", ".jpeg", ".png", ".webp", ".gif")):
return "image"
if _is_video(media_type) or lowered.endswith((".mp4", ".mov", ".mkv", ".avi", ".webm")):
return "video"
return "file"
def _collect_urls(node: Any) -> list[str]:
urls: list[str] = []
seen_ids: set[int] = set()
def walk(value: Any) -> None:
if value is None:
return
obj_id = id(value)
if obj_id in seen_ids:
return
seen_ids.add(obj_id)
if isinstance(value, str):
if value.startswith("http://") or value.startswith("https://"):
urls.append(value.strip())
return
if isinstance(value, (list, tuple, set)):
for item in value:
walk(item)
return
if isinstance(value, dict):
for nested in value.values():
walk(nested)
return
if hasattr(value, "__dict__"):
walk(vars(value))
walk(node)
return list(dict.fromkeys(urls))
def _extract_media_urls(message: Any) -> tuple[list[str], list[str], list[str], list[str]]: def _extract_media_urls(message: Any) -> tuple[list[str], list[str], list[str], list[str]]:
image_urls: list[str] = [] image_urls: list[str] = []
video_urls: list[str] = [] video_urls: list[str] = []
@@ -47,6 +105,7 @@ def _extract_media_urls(message: Any) -> tuple[list[str], list[str], list[str],
for item in raw_attachments: for item in raw_attachments:
data = _as_dict(item) data = _as_dict(item)
media_type = _stringify(data.get("type") or data.get("media_type") or data.get("kind")) media_type = _stringify(data.get("type") or data.get("media_type") or data.get("kind"))
is_forward_like = _is_forward_like(data)
url = _stringify( url = _stringify(
data.get("base_url") data.get("base_url")
or data.get("url") or data.get("url")
@@ -69,6 +128,23 @@ def _extract_media_urls(message: Any) -> tuple[list[str], list[str], list[str],
media_type = _stringify(nested_data.get("type") or nested_data.get("media_type")) media_type = _stringify(nested_data.get("type") or nested_data.get("media_type"))
if not url: if not url:
nested_urls = _collect_urls(item)
for nested_url in nested_urls:
kind = _classify_url(nested_url, media_type)
if kind == "image":
image_urls.append(nested_url)
elif kind == "video":
video_urls.append(nested_url)
else:
file_urls.append(nested_url)
if nested_urls:
continue
if not url:
if is_forward_like:
# Forward-пакет может не содержать прямого URL в верхнем уровне;
# текст/медиа достанем рекурсивно в других этапах.
continue
kind = media_type or _stringify(type(item).__name__) or "unknown" kind = media_type or _stringify(type(item).__name__) or "unknown"
unknown_attachments.append(kind) unknown_attachments.append(kind)
continue continue
+2
View File
@@ -11,6 +11,8 @@ class ParsedMessage:
image_urls: list[str] = field(default_factory=list) image_urls: list[str] = field(default_factory=list)
video_urls: list[str] = field(default_factory=list) video_urls: list[str] = field(default_factory=list)
file_urls: list[str] = field(default_factory=list) file_urls: list[str] = field(default_factory=list)
# URL -> исходное имя файла (если удалось определить в MAX).
file_names_by_url: dict[str, str] = field(default_factory=dict)
unknown_attachments: list[str] = field(default_factory=list) unknown_attachments: list[str] = field(default_factory=list)
# Ответ в MAX: Message.link указывает на исходное сообщение (тред). # Ответ в MAX: Message.link указывает на исходное сообщение (тред).
reply_to_max_message_id: str | None = None reply_to_max_message_id: str | None = None
+105 -6
View File
@@ -1,5 +1,9 @@
import asyncio import asyncio
import json import json
import os
import pathlib
import tempfile
import urllib.parse
from typing import Any from typing import Any
import requests import requests
@@ -39,6 +43,7 @@ class TelegramClient:
self._timeout = timeout self._timeout = timeout
self._chat_title_to_id: dict[str, str] = {} self._chat_title_to_id: dict[str, str] = {}
self._me: dict[str, Any] | None = None self._me: dict[str, Any] | None = None
self._tmp_root: str | None = None
@property @property
def fallback_user_id(self) -> str: def fallback_user_id(self) -> str:
@@ -107,19 +112,58 @@ class TelegramClient:
self, self,
chat_id: str, chat_id: str,
document_url: str, document_url: str,
file_name: str | None = None,
caption: str | None = None, caption: str | None = None,
*, *,
reply_to_message_id: int | None = None, reply_to_message_id: int | None = None,
) -> dict[str, Any]: ) -> dict[str, Any]:
payload: dict[str, Any] = { # Telegram часто не может скачать URL, которые доступны только клиенту MAX.
"chat_id": chat_id, # Поэтому скачиваем сами во временный файл и отправляем как multipart upload.
"document": document_url, safe_name = self._sanitize_filename(file_name) if file_name else None
} tmp_path = await self._download_to_temp(document_url, preferred_filename=safe_name)
try:
return await self.send_document_file(
chat_id=chat_id,
file_path=tmp_path,
upload_filename=safe_name,
caption=caption,
reply_to_message_id=reply_to_message_id,
)
finally:
try:
os.remove(tmp_path)
except OSError:
pass
async def send_document_file(
self,
*,
chat_id: str,
file_path: str,
upload_filename: str | None = None,
caption: str | None = None,
reply_to_message_id: int | None = None,
) -> dict[str, Any]:
payload: dict[str, Any] = {"chat_id": str(chat_id)}
if caption: if caption:
payload["caption"] = caption payload["caption"] = caption
if reply_to_message_id is not None: if reply_to_message_id is not None:
payload["reply_to_message_id"] = reply_to_message_id payload["reply_to_message_id"] = str(int(reply_to_message_id))
return await self._request("sendDocument", payload)
filename = upload_filename or pathlib.Path(file_path).name
def _do_request() -> requests.Response:
with open(file_path, "rb") as f:
files = {"document": (filename, f)}
return requests.post(
f"{self._base_url}/sendDocument",
data=payload,
files=files,
timeout=self._timeout,
)
response = await asyncio.to_thread(_do_request)
return self._parse_response(method="sendDocument", response=response)
async def send_media_group( async def send_media_group(
self, self,
@@ -249,6 +293,9 @@ class TelegramClient:
return requests.post(url, json=payload, timeout=self._timeout) return requests.post(url, json=payload, timeout=self._timeout)
response = await asyncio.to_thread(_do_request) response = await asyncio.to_thread(_do_request)
return self._parse_response(method=method, response=response)
def _parse_response(self, *, method: str, response: requests.Response) -> dict[str, Any]:
if response.status_code >= 400: if response.status_code >= 400:
error_code: int | None = None error_code: int | None = None
description: str | None = None description: str | None = None
@@ -283,3 +330,55 @@ class TelegramClient:
parameters=data.get("parameters") if isinstance(data.get("parameters"), dict) else None, parameters=data.get("parameters") if isinstance(data.get("parameters"), dict) else None,
) )
return data return data
async def _download_to_temp(self, url: str, preferred_filename: str | None = None) -> str:
filename = preferred_filename or self._infer_filename_from_url(url) or "max-file"
tmp_dir = self._ensure_tmp_root()
fd, path = tempfile.mkstemp(prefix="max2tg_", suffix=f"_{filename}", dir=tmp_dir)
os.close(fd)
def _do_download() -> None:
with requests.get(url, stream=True, timeout=self._timeout) as r:
r.raise_for_status()
with open(path, "wb") as f:
for chunk in r.iter_content(chunk_size=1024 * 256):
if chunk:
f.write(chunk)
try:
await asyncio.to_thread(_do_download)
return path
except Exception:
try:
os.remove(path)
except OSError:
pass
raise
def _ensure_tmp_root(self) -> str:
if self._tmp_root and os.path.isdir(self._tmp_root):
return self._tmp_root
self._tmp_root = tempfile.mkdtemp(prefix="max2tg_")
return self._tmp_root
@staticmethod
def _infer_filename_from_url(url: str) -> str | None:
try:
parsed = urllib.parse.urlparse(url)
name = pathlib.Path(parsed.path).name
if name and name not in {"/", ".", ".."}:
# Windows-safe filename (и вообще безопаснее для FS).
bad = '<>:"/\\|?*'
cleaned = "".join("_" if ch in bad else ch for ch in name).strip().strip(".")
return cleaned or None
except Exception:
return None
return None
@staticmethod
def _sanitize_filename(name: str | None) -> str | None:
if not name:
return None
bad = '<>:"/\\|?*'
cleaned = "".join("_" if ch in bad else ch for ch in str(name)).strip().strip(".")
return cleaned or None
+473
View File
@@ -0,0 +1,473 @@
# Техническое задание: двунаправленный мост MAX <-> Telegram
## 1. Назначение системы
Система должна обеспечивать непрерывную двунаправленную синхронизацию сообщений между мессенджером MAX и Telegram:
- направление A: входящие сообщения из MAX пересылаются в Telegram;
- направление B: входящие сообщения из Telegram пересылаются в MAX;
- поддерживаются текст, фото, видео, документы/файлы и ответы (reply) в пределах доступных API;
- реализованы дедупликация, хранение связей между сообщениями в обоих направлениях и базовый health-monitoring;
- предусмотрен fallback-маршрут в Telegram при отсутствии явного соответствия чатов.
Система реализуется как один сервисный процесс (или контейнер), работающий постоянно.
## 2. Контекст и границы
### 2.1 Внешние зависимости
Обязательные внешние системы:
- API MAX (через клиентскую библиотеку или собственный API-клиент);
- Telegram Bot API (через клиентскую библиотеку);
- SQLite (или совместимое хранилище) для локального состояния;
- переменные окружения для конфигурации.
### 2.2 Что входит в систему
- запуск и аутентификация клиента MAX;
- polling Telegram `getUpdates` (без webhook);
- обработка входящих событий в обоих направлениях;
- маршрутизация по названию чата и по явным биндам;
- хранение mapping/дедупликации/маршрутов;
- HTTP health endpoints (`/livez`, `/healthz`).
### 2.3 Что НЕ входит в систему
- UI/панель администрирования;
- сложная очередь сообщений (Kafka/RabbitMQ и т.п.);
- гарантированная exactly-once доставка между платформами;
- миграции БД с версионированием (в базовой реализации только `CREATE TABLE IF NOT EXISTS`);
- хранение медиа в собственной файловой инфраструктуре.
## 3. Функциональные требования
## 3.1 MAX -> Telegram
При получении сообщения из MAX система должна:
1. Распарсить сообщение:
- `message_id`, `chat_id`, `sender_name`, `chat_name`, текст;
- список `image_urls`, `video_urls`, `file_urls`;
- список неизвестных вложений;
- данные reply-контекста (`reply_to_max_message_id`, preview).
2. Обогатить данные через API MAX (best effort):
- попытаться получить человекочитаемое имя отправителя;
- попытаться получить реальное название чата;
- извлечь дополнительные вложения через типизированные attach-объекты.
3. Проверить дедупликацию по паре `(message_id, chat_id)`; дубликаты не отправлять.
4. Выбрать Telegram-чат:
- сначала через явный route (bind) по нормализованному названию MAX-чата;
- затем через кэш чатов Telegram по совпадению заголовка;
- если не найдено — отправка в `fallback_user_id`.
5. Сформировать текст:
- если целевой Telegram-чат найден (не fallback): `"{sender}:\n{text}"`;
- если fallback: `"{sender} / {chat}:\n{text}"`;
- при неизвестных вложениях добавить уведомление в конец.
6. Обработать reply:
- попытаться найти Telegram `reply_to_message_id` через mapping;
- если mapping не найден — добавить текстовый контекст ответа.
7. Отправить контент:
- при `image+video > 1` — отправить единым альбомом (`sendMediaGroup`);
- иначе отправить текст/медиа/файлы поштучно;
- при полностью пустом payload отправить служебный fallback-текст.
8. Сохранить mapping отправленных сообщений.
9. Пометить сообщение как forwarded.
10. При ошибке основного пути сделать аварийное best-effort уведомление в fallback-чат Telegram.
## 3.2 Telegram -> MAX
Система должна запускать единственный polling-цикл `getUpdates` и:
1. При старте получить `bot_id` через `getMe`.
2. Обновить локальный кэш чатов MAX (`title -> id`) на основе доступного списка.
3. В цикле получать updates с `offset` и `allowed_updates=["message","channel_post"]`.
4. Для каждого сообщения:
- отбросить неподдерживаемые типы;
- отбросить сообщения бота (защита от петель);
- обработать команду `/bind_max` в приоритетном порядке;
- обработать служебные команды управления MAX (только личка + только fallback user);
- для обычных сообщений найти чат MAX по нормализованному названию Telegram-чата;
- если чат MAX не найден — не пересылать, логировать ошибку.
5. Обработать `media_group_id`:
- буферизовать элементы альбома по ключу `(telegram_chat_id, media_group_id)`;
- после grace-паузы (около 1.2 сек) отправить одним сообщением в MAX с несколькими вложениями.
6. Для одиночных сообщений:
- сформировать текст `"{sender}:\n{text_or_caption}"`;
- для фото/видео получить URL через `getFile` и вложить в MAX как native media attach;
- для document/audio/voice/animation/sticker/video_note добавить URL-список в текстовый блок;
- при reply попытаться найти соответствующее сообщение MAX через mapping.
7. После успешной отправки:
- поставить реакцию на Telegram-сообщение (emoji, best effort);
- сохранить mapping `(telegram chat/message -> max chat/message)`.
## 3.3 Управляющие команды в Telegram
Команды обрабатываются только в приватном чате с ботом и только от пользователя `fallback_user_id`, кроме `/bind_max` (она может работать в целевом чате).
Поддерживаемые команды:
- `/help` — список команд;
- `/list` — список активных чатов MAX;
- `/join <LINK>` — вступление в группу/канал MAX по ссылке;
- `/leave <НАЗВАНИЕ>` — выход из канала MAX (только если тип чата определен как канал);
- `/last_messages <НАЗВАНИЕ>` — последние 10 сообщений;
- `/bind_max <точное название MAX-чата>` — привязка текущего Telegram-чата к MAX-чату.
Требования:
- у команд должны быть понятные текстовые ответы;
- ошибки внешних API должны возвращаться в ответе, не падая процессом;
- `/bind_max` должна валидировать существование MAX-чата перед сохранением маршрута.
## 4. Нефункциональные требования
- **Надежность:** сервис работает бесконечно, при ошибках polling применяет backoff.
- **Идемпотентность (частичная):** дедупликация MAX->Telegram через БД.
- **Наблюдаемость:** структурированные логи и health endpoint.
- **Портируемость:** реализация возможна на любом языке при соблюдении контрактов.
- **Производительность:** обработка событий в near real-time, без тяжелых batch-процессов.
- **Отказоустойчивость:** best-effort при частичных отказах API/медиа.
## 5. Конфигурация (env contract)
Обязательные параметры:
- `MAX_PHONE` — телефон аккаунта MAX (для авторизации/сессии);
- `MAX_WORK_DIR` — рабочая директория клиента MAX (по умолчанию `cache`);
- `TELEGRAM_BOT_TOKEN` — токен Telegram-бота;
- `TELEGRAM_FALLBACK_USER_ID` — Telegram user/chat id для fallback;
- `SQLITE_PATH` — путь до SQLite БД (по умолчанию `${MAX_WORK_DIR}/max2telegram.db`).
Дополнительно:
- `TZ` — таймзона окружения;
- `PYTHONUNBUFFERED` (или аналог) — политика буферизации логов.
Валидация:
- обязательные переменные должны проверяться при старте;
- при отсутствии обязательной переменной процесс завершает запуск с явной ошибкой.
## 6. Архитектура
### 6.1 Компоненты
1. **Bootstrap / Main**
- загружает конфиг;
- создает клиентов MAX и Telegram;
- инициализирует storage;
- запускает health server;
- подключает обработчики входящих сообщений MAX;
- запускает Telegram->MAX poller как фоновую задачу.
2. **MAX->Telegram Bridge**
- парсинг + enrich входящих MAX-сообщений;
- маршрутизация в Telegram;
- отправка текст/медиа/документы;
- обработка миграции Telegram chat id;
- сохранение mapping и dedup.
3. **Telegram API Client**
- HTTP-обертка над Bot API;
- методы отправки всех типов контента;
- `getUpdates` + кэширование известных чатов;
- `getFile` для медиа URL;
- унифицированная модель ошибок с извлечением `migrate_to_chat_id`.
4. **Telegram->MAX Bridge**
- polling updates;
- фильтрация собственных сообщений;
- обработка команд;
- преобразование Telegram payload -> MAX message/attachments;
- буферизация media group;
- реакция и mapping.
5. **Parser MAX Message**
- универсальный best-effort разбор разнородных форматов вложений;
- классификация URL по типам;
- извлечение reply-контекста.
6. **Storage**
- таблица дедупликации;
- таблица двустороннего mapping;
- таблица явных маршрутов (binds).
7. **Health subsystem**
- хранит отметки времени последнего успеха/ошибки по MAX и Telegram;
- формирует snapshot;
- HTTP endpoint отдает liveness/readiness.
8. **Auth utility**
- отдельный скрипт первичной авторизации MAX;
- сохраняет сессию в рабочем каталоге.
### 6.2 Логическая схема взаимодействия
- MAX event -> Parser -> MAX->TG Bridge -> Telegram API -> Storage update.
- Telegram update -> TG->MAX Bridge -> MAX API -> Storage update.
- TG->MAX Bridge единолично вызывает `getUpdates`, одновременно наполняя кэш чатов Telegram.
- Storage используется обоими мостами как разделяемый слой состояния.
- Health обновляется из polling-циклов и MAX событий.
## 7. Модель данных и БД
Используется SQLite (или эквивалент в другой СУБД).
### 7.1 Таблица `forwarded_messages`
Назначение: дедупликация MAX->Telegram.
Поля:
- `message_id TEXT NOT NULL`
- `chat_id TEXT NOT NULL`
- `forwarded_at DATETIME DEFAULT CURRENT_TIMESTAMP`
Ключ:
- `PRIMARY KEY (message_id, chat_id)`
### 7.2 Таблица `message_mapping`
Назначение: двусторонняя связка сообщений для reply и трассировки.
Поля:
- `telegram_chat_id TEXT NOT NULL`
- `telegram_message_id TEXT NOT NULL`
- `max_chat_id TEXT NOT NULL`
- `max_message_id TEXT NOT NULL`
- `media_group_id TEXT NULL`
- `created_at DATETIME DEFAULT CURRENT_TIMESTAMP`
Ключ:
- `PRIMARY KEY (telegram_chat_id, telegram_message_id)`
Индексы:
- `(max_chat_id, max_message_id)` — поиск Telegram-сообщения по MAX;
- `(telegram_chat_id, media_group_id)` — групповые операции альбомов.
### 7.3 Таблица `chat_routes`
Назначение: явные маршруты MAX chat title -> Telegram chat id.
Поля:
- `max_chat_title_norm TEXT PRIMARY KEY`
- `telegram_chat_id TEXT NOT NULL`
- `telegram_chat_title TEXT NULL`
- `created_at DATETIME DEFAULT CURRENT_TIMESTAMP`
Нормализация ключа:
- trim + casefold/lower (без учета регистра).
## 8. Алгоритмы и правила
## 8.1 Нормализация названий чатов
Во всех маршрутизирующих сравнениях:
- удалить крайние пробелы;
- привести к регистронезависимой форме (`casefold`/`lower`);
- сравнивать только в нормализованном виде.
## 8.2 Политика маршрутизации MAX -> Telegram
Порядок выбора:
1. `chat_routes` по нормализованному MAX title;
2. локальный кэш Telegram title->id (наполняется из `getUpdates`);
3. fallback user/chat id.
## 8.3 Политика медиа
- MAX->Telegram:
- если суммарно фото+видео больше одного, использовать album API;
- документы отправлять отдельными сообщениями;
- caption добавлять только к первому элементу/первому отправляемому сообщению.
- Telegram->MAX:
- фото/видео отправлять как native attachments;
- прочие типы файлов прикладывать ссылками в тексте.
## 8.4 Reply-семантика
- Для Telegram->MAX:
- если Telegram message является reply, искать соответствующий `max_message_id` в mapping;
- если найден, отправлять `reply_to` в MAX.
- Для MAX->Telegram:
- если MAX message является reply, искать `telegram_message_id` в mapping;
- если найден, отправлять `reply_to_message_id`;
- если не найден, добавлять текстовую пометку с превью исходного сообщения.
## 8.5 Обработка migration в Telegram
Если Telegram API возвращает ошибку с `parameters.migrate_to_chat_id`:
1. обновить маршрут в `chat_routes` на новый chat id;
2. повторить отправку в новый chat id;
3. считать повтор успешным итоговым результатом.
## 8.6 Дедупликация
- применяется для MAX->Telegram по ключу `(max_message_id, max_chat_id)`;
- после успешной отправки обязательно mark-forwarded;
- на старте/рестарте состояния берутся из БД.
## 8.7 Buffering media group (Telegram)
- ключ буфера: `(telegram_chat_id, media_group_id)`;
- каждое сообщение альбома копится в списке;
- по истечении grace-периода группа отправляется одним вызовом в MAX;
- после отправки буфер очищается.
## 9. API-контракты внутренних модулей
Ниже абстрактные контракты, независимые от языка:
- `Settings load_settings()`
- читает env;
- валидирует обязательные поля;
- возвращает immutable-конфигурацию.
- `ParsedMessage parse_message(MaxMessage msg)`
- best-effort преобразование сырого MAX-сообщения в каноническую структуру.
- `BridgeStorage`
- `was_forwarded(message_id, chat_id) -> bool`
- `mark_forwarded(message_id, chat_id)`
- `save_mapping(telegram_chat_id, telegram_message_id, max_chat_id, max_message_id, media_group_id?)`
- `get_max_message_id_for_telegram(telegram_chat_id, telegram_message_id) -> str?`
- `get_telegram_message_id_for_max(telegram_chat_id, max_chat_id, max_message_id) -> str?`
- `set_chat_route(max_chat_title_norm, telegram_chat_id, telegram_chat_title?)`
- `get_chat_route(max_chat_title_norm) -> str?`
- `TelegramClient`
- `resolve_target_chat_id(max_chat_name) -> (chat_id, matched_by_title)`
- `send_text/send_photo/send_video/send_document/send_media_group(...)`
- `get_updates(offset, timeout, limit) -> updates[]`
- `get_file_url(file_id) -> url`
- `add_reaction(chat_id, message_id, emoji)`
- `MaxToTelegramBridge.forward_message(max_message)`
- `TelegramToMaxBridge.start()`
- `handle_control_command(message, max_client, telegram) -> str?`
## 10. Health-check модель
Должны храниться timestamps:
- telegram: `last_ok`, `last_error`;
- max: `last_ok`, `last_error`, `last_event`;
- `started_at`.
Параметр:
- `unhealthy_after_sec` (по умолчанию ~15 минут).
Правила:
- `telegram_healthy = last_ok exists && now-last_ok <= unhealthy_after_sec`;
- `max_healthy = last_ok exists && now-last_ok <= unhealthy_after_sec`;
- `overall_healthy = telegram_healthy && max_healthy`.
HTTP:
- `GET /livez` (и `/live`, `/`) -> 200, `{status:"live", uptime_sec}`.
- `GET /healthz` (и `/health`) -> 200 или 503, детальный JSON по компонентам.
## 11. Поведение при ошибках и устойчивость
- Любая ошибка обработки отдельного сообщения не должна останавливать сервис.
- Polling Telegram при ошибке уходит в backoff (например 10 секунд).
- Ошибка реакции в Telegram не влияет на основную доставку.
- Ошибка аварийного уведомления логируется, но не роняет процесс.
- Проблемы с разрешением URL медиа обрабатываются best-effort:
- что удалось достать — отправляется;
- что не удалось — отражается в тексте/unknown notices.
## 12. Логирование и диагностика
Обязательные события логов:
- старт/остановка компонентов;
- маршрутизация сообщений;
- обнаружение дубликатов;
- ошибки API и stacktrace;
- успешная пересылка с количеством вложений;
- операции bind и migration chat id.
Рекомендуемый формат:
- timestamp + level + logger + message.
## 13. Сценарии запуска и деплой
## 13.1 Локальный запуск
1. Подготовить `.env`.
2. Установить зависимости.
3. Один раз пройти auth MAX (интерактивно), сохранить сессию.
4. Запустить основной процесс.
## 13.2 Контейнерный запуск
- контейнер должен включать runtime + зависимости;
- каталог `cache` должен быть volume для сохранения сессии и SQLite;
- должен быть healthcheck через `GET /healthz`.
## 13.3 CI/CD (рекомендованно)
- сборка Docker image при push в основные ветки;
- публикация в registry с тегами для dev/release/latest;
- кэширование слоев сборки.
## 14. Требования к переносимой реализации (на любом языке)
Чтобы воссоздать проект в другом языке, необходимо сохранить:
1. Два независимых, но согласованных канала обработки:
- event-driven для MAX сообщений;
- polling-loop для Telegram updates.
2. Единое постоянное хранилище с тремя сущностями:
- dedup;
- mapping;
- routes.
3. Идентичные правила нормализации названий чатов и выбора маршрута.
4. Reply-механику с fallback-текстом при отсутствии mapping.
5. Политику media group и порядок отправки вложений.
6. Обработку Telegram migration (`migrate_to_chat_id`) с обновлением маршрута.
7. Health-модель с независимыми метками MAX/Telegram.
8. Ограничение: только один consumer `getUpdates` на экземпляр бота.
## 15. Acceptance criteria
Система считается реализованной, если:
1. Текст/фото/видео/файлы корректно ходят в обе стороны.
2. MAX->Telegram не дублирует уже пересланные сообщения после рестарта.
3. Reply в обе стороны сохраняется, если mapping существует.
4. При отсутствии Telegram-совпадения сообщение уходит в fallback.
5. `/bind_max` меняет маршрут и влияет на последующие MAX->Telegram сообщения.
6. Telegram media group приходит в MAX как одно сообщение с множеством вложений.
7. `/healthz` возвращает 200 при рабочем MAX+Telegram и 503 при деградации.
8. Ошибки отдельных сообщений не приводят к остановке процесса.
## 16. Известные ограничения текущей логики
- Автопоиск Telegram-чата зависит от того, что чат уже встречался в `getUpdates`.
- Для нестандартных вложений возможна частичная деградация в plain text + ссылки.
- Дедупликация реализована только для потока MAX->Telegram.
- Конкурентный доступ к SQLite через множество соединений допустим для небольших нагрузок, но для high-load может потребоваться иной storage backend.
## 17. Рекомендации для расширения (необязательно)
- добавить миграции схемы БД;
- добавить retry policy с классификацией transient/permanent ошибок;
- добавить метрики (Prometheus/OpenTelemetry);
- добавить интеграционные тесты с моками API MAX/Telegram;
- добавить персистентный offset Telegram (если нужен recovery без повторов между рестартами).