16 Commits
Author SHA1 Message Date
kislovdm 64fe9c466e маршруты
Docker Hub / build-and-push (push) Failing after 15s
2026-04-23 19:56:20 +03:00
kislovdm 8bba0d9dd5 tech spec 2026-04-22 14:52:03 +03:00
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
kislovdm e40297787e temporary skip max
Docker Hub / build-and-push (push) Failing after 13s
2026-04-20 19:03:44 +03:00
kislovdm dcf4441340 work with files 2026-04-20 18:55:01 +03:00
kislovdm 1ef5cf4e38 ИСправление списка каналов
Docker Hub / build-and-push (push) Failing after 14s
2026-04-15 14:30:45 +03:00
kislovdm 1fd5e4ea39 после миграции группы в супергруппу необходимо обновить её id
Docker Hub / build-and-push (push) Failing after 14s
2026-04-14 11:13:21 +03:00
kislovdm bcf9366cf6 remove sha builds 2026-04-11 11:00:32 +03:00
kislovdm 102009369a clear caches in image 2026-04-11 10:58:08 +03:00
kislovdm 7a3c7317ea fix actions 2026-04-11 10:55:07 +03:00
11 changed files with 1259 additions and 65 deletions
+6 -3
View File
@@ -6,6 +6,10 @@
# #
# Опционально: Variables → DOCKERHUB_IMAGE (например org/max2telegram), если имя образа # Опционально: Variables → DOCKERHUB_IMAGE (например org/max2telegram), если имя образа
# отличается от <DOCKERHUB_USERNAME>/max2telegram. # отличается от <DOCKERHUB_USERNAME>/max2telegram.
#
# Теги образа:
# push в master → dev (+ sha)
# push тега v* → {{version}}, {{major}}.{{minor}}, latest (+ sha)
name: Docker Hub name: Docker Hub
@@ -57,11 +61,10 @@ jobs:
with: with:
images: ${{ steps.image.outputs.name }} images: ${{ steps.image.outputs.name }}
tags: | tags: |
type=ref,event=branch type=raw,value=dev,enable=${{ github.ref == 'refs/heads/master' }}
type=semver,pattern={{version}} type=semver,pattern={{version}}
type=semver,pattern={{major}}.{{minor}} type=semver,pattern={{major}}.{{minor}}
type=sha,prefix= type=raw,value=latest,enable=${{ startsWith(github.ref, 'refs/tags/') }}
type=raw,value=latest,enable={{is_default_branch}}
- name: Build and push - name: Build and push
uses: docker/build-push-action@v6 uses: docker/build-push-action@v6
+3 -1
View File
@@ -2,7 +2,9 @@ FROM python:3.12-slim
WORKDIR /app WORKDIR /app
RUN apt-get update; apt-get install -y git RUN apt-get update \
&& apt-get install -y --no-install-recommends git \
&& rm -rf /var/lib/apt/lists/*
COPY ./src/requirements.txt . COPY ./src/requirements.txt .
RUN pip install --no-cache-dir -r ./requirements.txt RUN pip install --no-cache-dir -r ./requirements.txt
+366 -33
View File
@@ -1,12 +1,13 @@
import logging import logging
from collections.abc import Awaitable, Callable
from typing import Any 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 TelegramClient from telegram_api import TelegramApiError, TelegramClient
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@@ -18,9 +19,9 @@ class MaxToTelegramBridge:
self._storage = storage self._storage = storage
async def forward_message(self, max_message: Any) -> None: async def forward_message(self, max_message: Any) -> None:
if self._is_self_message(max_message): #if self._is_self_message(max_message):
logger.debug("Skip self message %s/%s", getattr(max_message, "chat_id", "?"), getattr(max_message, "id", "?")) # logger.debug("Skip self message %s/%s", getattr(max_message, "chat_id", "?"), getattr(max_message, "id", "?"))
return # return
parsed = parse_message(max_message) parsed = parse_message(max_message)
parsed = await self._enrich_from_max(max_message, parsed) parsed = await self._enrich_from_max(max_message, parsed)
@@ -33,11 +34,6 @@ class MaxToTelegramBridge:
# - если найден целевой Telegram-чат (не fallback): "Ирина:\n<текст>" # - если найден целевой Telegram-чат (не fallback): "Ирина:\n<текст>"
# - если fallback: "Ирина / Свободный микрофон:\n<текст>" # - если fallback: "Ирина / Свободный микрофон:\n<текст>"
# Решение о том, включать ли название чата, принимаем после определения маршрута. # Решение о том, включать ли название чата, принимаем после определения маршрута.
has_any_payload = bool((parsed.text or "").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
normalized = parsed.chat_name.strip().casefold() normalized = parsed.chat_name.strip().casefold()
routed = self._storage.get_chat_route(max_chat_title_norm=normalized) routed = self._storage.get_chat_route(max_chat_title_norm=normalized)
if routed: if routed:
@@ -63,6 +59,7 @@ class MaxToTelegramBridge:
text=parsed.text, text=parsed.text,
include_chat_name=is_fallback, include_chat_name=is_fallback,
) )
text = self._append_unknown_attachment_notice(parsed=parsed, text=text)
reply_telegram_mid = self._resolve_telegram_reply_to( reply_telegram_mid = self._resolve_telegram_reply_to(
telegram_chat_id=str(target_chat_id), telegram_chat_id=str(target_chat_id),
@@ -72,20 +69,25 @@ class MaxToTelegramBridge:
if parsed.reply_to_max_message_id and reply_telegram_mid is None: if parsed.reply_to_max_message_id and reply_telegram_mid is None:
text = self._prepend_max_reply_context(parsed, text) text = self._prepend_max_reply_context(parsed, text)
has_any_payload = bool(text.strip()) or bool(parsed.image_urls) or bool(parsed.video_urls) has_any_payload = bool(text.strip()) or bool(parsed.image_urls) or bool(parsed.video_urls) or bool(parsed.file_urls)
if not has_any_payload: if not has_any_payload:
logger.debug("Skip empty message %s/%s", parsed.chat_id, parsed.message_id) text = self._build_fallback_unknown_notice(parsed)
return
total_media = len(parsed.image_urls) + len(parsed.video_urls) total_media = len(parsed.image_urls) + len(parsed.video_urls)
sent_any = False
if total_media > 1: if total_media > 1:
# Отправляем одним альбомом в Telegram (единое сообщение). # Отправляем одним альбомом в Telegram (единое сообщение).
sent_messages = await self._telegram.send_media_group( target_chat_id, sent_messages = await self._send_with_migration_retry(
chat_id=target_chat_id, target_chat_id=target_chat_id,
max_chat_title_norm=normalized,
max_chat_title=parsed.chat_name,
send_action=lambda chat_id: self._telegram.send_media_group(
chat_id=chat_id,
image_urls=parsed.image_urls, image_urls=parsed.image_urls,
video_urls=parsed.video_urls, video_urls=parsed.video_urls,
caption=text, caption=text,
reply_to_message_id=reply_telegram_mid, reply_to_message_id=reply_telegram_mid,
),
) )
for sent in sent_messages: for sent in sent_messages:
mid = sent.get("message_id") mid = sent.get("message_id")
@@ -97,6 +99,7 @@ class MaxToTelegramBridge:
max_chat_id=str(parsed.chat_id), max_chat_id=str(parsed.chat_id),
max_message_id=str(parsed.message_id), max_message_id=str(parsed.message_id),
) )
sent_any = True
self._storage.mark_forwarded(parsed.message_id, parsed.chat_id) self._storage.mark_forwarded(parsed.message_id, parsed.chat_id)
logger.info( logger.info(
"Forwarded media group %s/%s (images=%s, videos=%s)", "Forwarded media group %s/%s (images=%s, videos=%s)",
@@ -107,10 +110,15 @@ class MaxToTelegramBridge:
) )
return return
sent_any = False should_send_plain_text = total_media == 0 and not parsed.file_urls and bool(text.strip())
if parsed.text.strip() and total_media == 0: if should_send_plain_text:
sent = await self._telegram.send_text( target_chat_id, sent = await self._send_with_migration_retry(
target_chat_id, text, reply_to_message_id=reply_telegram_mid target_chat_id=target_chat_id,
max_chat_title_norm=normalized,
max_chat_title=parsed.chat_name,
send_action=lambda chat_id: self._telegram.send_text(
chat_id, 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
if mid is not None: if mid is not None:
@@ -124,11 +132,16 @@ class MaxToTelegramBridge:
for index, image_url in enumerate(parsed.image_urls): for index, image_url in enumerate(parsed.image_urls):
caption = text if not sent_any and index == 0 else None caption = text if not sent_any and index == 0 else None
sent = await self._telegram.send_photo( target_chat_id, sent = await self._send_with_migration_retry(
target_chat_id, target_chat_id=target_chat_id,
max_chat_title_norm=normalized,
max_chat_title=parsed.chat_name,
send_action=lambda chat_id: self._telegram.send_photo(
chat_id,
image_url, image_url,
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,
),
) )
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
if mid is not None: if mid is not None:
@@ -142,11 +155,62 @@ class MaxToTelegramBridge:
for index, video_url in enumerate(parsed.video_urls): for index, video_url in enumerate(parsed.video_urls):
caption = text if not sent_any and index == 0 else None caption = text if not sent_any and index == 0 else None
sent = await self._telegram.send_video( target_chat_id, sent = await self._send_with_migration_retry(
target_chat_id, target_chat_id=target_chat_id,
max_chat_title_norm=normalized,
max_chat_title=parsed.chat_name,
send_action=lambda chat_id: self._telegram.send_video(
chat_id,
video_url, video_url,
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,
),
)
mid = sent.get("result", {}).get("message_id") if isinstance(sent.get("result"), dict) else None
if mid is not None:
self._storage.save_mapping(
telegram_chat_id=str(target_chat_id),
telegram_message_id=str(mid),
max_chat_id=str(parsed.chat_id),
max_message_id=str(parsed.message_id),
)
sent_any = True
for index, file_url in enumerate(parsed.file_urls):
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=target_chat_id,
max_chat_title_norm=normalized,
max_chat_title=parsed.chat_name,
send_action=lambda chat_id: self._telegram.send_document(
chat_id,
file_url,
file_name=file_name,
caption=caption,
reply_to_message_id=reply_telegram_mid if not sent_any and index == 0 else None,
),
)
mid = sent.get("result", {}).get("message_id") if isinstance(sent.get("result"), dict) else None
if mid is not None:
self._storage.save_mapping(
telegram_chat_id=str(target_chat_id),
telegram_message_id=str(mid),
max_chat_id=str(parsed.chat_id),
max_message_id=str(parsed.message_id),
)
sent_any = True
if not sent_any:
# Последняя страховка: гарантируем уведомление в 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=target_chat_id,
max_chat_title_norm=normalized,
max_chat_title=parsed.chat_name,
send_action=lambda chat_id: self._telegram.send_text(
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
if mid is not None: if mid is not None:
@@ -160,13 +224,64 @@ class MaxToTelegramBridge:
self._storage.mark_forwarded(parsed.message_id, parsed.chat_id) self._storage.mark_forwarded(parsed.message_id, parsed.chat_id)
logger.info( logger.info(
"Forwarded message %s/%s (images=%s, videos=%s)", "Forwarded message %s/%s (images=%s, videos=%s, files=%s, unknown=%s)",
parsed.chat_id, parsed.chat_id,
parsed.message_id, parsed.message_id,
len(parsed.image_urls), len(parsed.image_urls),
len(parsed.video_urls), len(parsed.video_urls),
len(parsed.file_urls),
len(parsed.unknown_attachments),
) )
async def notify_delivery_failure(self, max_message: Any, error: Exception) -> None:
"""Best-effort аварийное уведомление, если основной форвардинг упал."""
try:
parsed = parse_message(max_message)
parsed = await self._enrich_from_max(max_message, parsed)
body = self._build_fallback_unknown_notice(parsed)
body = f"{body}\n\n[bridge-error] {type(error).__name__}: {error}"
except Exception:
body = f"[!] Сообщение из MAX не доставлено в Telegram из-за ошибки bridge: {type(error).__name__}: {error}"
fallback_chat_id = self._telegram.fallback_user_id
if not fallback_chat_id:
logger.error("Cannot send emergency notice: Telegram fallback user id is empty")
return
try:
await self._telegram.send_text(chat_id=fallback_chat_id, text=body)
logger.warning("Sent emergency notice to Telegram fallback chat %s", fallback_chat_id)
except Exception:
logger.exception("Cannot send emergency notice to Telegram")
async def _send_with_migration_retry(
self,
*,
target_chat_id: str,
max_chat_title_norm: str,
max_chat_title: str,
send_action: Callable[[str], Awaitable[Any]],
) -> tuple[str, Any]:
try:
sent = await send_action(target_chat_id)
return target_chat_id, sent
except TelegramApiError as exc:
migrated_chat_id = exc.migrate_to_chat_id
if not migrated_chat_id or migrated_chat_id == str(target_chat_id):
raise
logger.warning(
"Telegram chat %s upgraded to %s for MAX chat '%s'; update route and retry",
target_chat_id,
migrated_chat_id,
max_chat_title,
)
self._storage.set_chat_route(
max_chat_title_norm=max_chat_title_norm,
telegram_chat_id=migrated_chat_id,
telegram_chat_title=max_chat_title,
)
sent = await send_action(migrated_chat_id)
return migrated_chat_id, sent
def _resolve_telegram_reply_to( def _resolve_telegram_reply_to(
self, *, telegram_chat_id: str, max_chat_id: str, parsed: ParsedMessage self, *, telegram_chat_id: str, max_chat_id: str, parsed: ParsedMessage
) -> int | None: ) -> int | None:
@@ -227,27 +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
# Убираем дубли URL, если парсер и enrich нашли одинаковые вложения. if isinstance(attach, FileAttach):
parsed.image_urls = list(dict.fromkeys(parsed.image_urls)) resolved = await self._resolve_file_attach_url(
parsed.video_urls = list(dict.fromkeys(parsed.video_urls)) message_chat_id=message_chat_id,
return parsed 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:
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)
if urls:
parsed.file_urls.extend(urls)
continue
parsed.unknown_attachments.append(type(attach).__name__)
async def _resolve_file_attach_url(self, *, message_chat_id: Any, message_id: Any, attach: FileAttach) -> str | None:
file_id = getattr(attach, "file_id", None)
if file_id is None or message_id is None:
return None
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)
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)
@@ -295,3 +555,76 @@ class MaxToTelegramBridge:
walk(attach) walk(attach)
return list(dict.fromkeys(urls)) return list(dict.fromkeys(urls))
def _extract_any_urls(self, 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)
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))
@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
def _append_unknown_attachment_notice(*, parsed: ParsedMessage, text: str) -> str:
if not parsed.unknown_attachments:
return text
unknown_preview = ", ".join(parsed.unknown_attachments[:5])
suffix = f"\n\n[!] Неизвестный тип вложения из MAX: {unknown_preview}"
return f"{text}{suffix}" if text else suffix.strip()
@staticmethod
def _build_fallback_unknown_notice(parsed: ParsedMessage) -> str:
base = MaxToTelegramBridge._format_caption(
sender_name=parsed.sender_name,
chat_name=parsed.chat_name,
text=parsed.text,
include_chat_name=True,
)
unknown = ", ".join(parsed.unknown_attachments[:5]) if parsed.unknown_attachments else "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"}
+2 -1
View File
@@ -72,8 +72,9 @@ def main() -> None:
health.mark_max_event() health.mark_max_event()
try: try:
await bridge.forward_message(message) await bridge.forward_message(message)
except Exception: except Exception as exc:
logger.exception("Failed to forward Max message") logger.exception("Failed to forward Max message")
await bridge.notify_delivery_failure(message, exc)
asyncio.run(max_client.start()) asyncio.run(max_client.start())
+87 -3
View File
@@ -36,15 +36,76 @@ 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 _extract_media_urls(message: Any) -> tuple[list[str], list[str]]: 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]]:
image_urls: list[str] = [] image_urls: list[str] = []
video_urls: list[str] = [] video_urls: list[str] = []
file_urls: list[str] = []
unknown_attachments: list[str] = []
# В PyMax рабочее поле для вложений обычно называется attaches. # В PyMax рабочее поле для вложений обычно называется attaches.
raw_attachments = _get_attr(message, ["attaches", "attachments", "media", "files"], default=[]) or [] raw_attachments = _get_attr(message, ["attaches", "attachments", "media", "files"], default=[]) or []
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")
@@ -67,14 +128,35 @@ def _extract_media_urls(message: Any) -> tuple[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"
unknown_attachments.append(kind)
continue continue
if _is_image(media_type): if _is_image(media_type):
image_urls.append(url) image_urls.append(url)
elif _is_video(media_type): elif _is_video(media_type):
video_urls.append(url) video_urls.append(url)
else:
file_urls.append(url)
return image_urls, video_urls return image_urls, video_urls, file_urls, unknown_attachments
def _extract_max_reply(message: Any) -> tuple[str | None, str | None]: def _extract_max_reply(message: Any) -> tuple[str | None, str | None]:
@@ -115,7 +197,7 @@ def parse_message(message: Any) -> ParsedMessage:
chat_id = _stringify(_get_attr(message, ["chat_id", "dialog_id", "peer_id"])) or "unknown-chat" chat_id = _stringify(_get_attr(message, ["chat_id", "dialog_id", "peer_id"])) or "unknown-chat"
text = _stringify(_get_attr(message, ["text", "message", "body"])) text = _stringify(_get_attr(message, ["text", "message", "body"]))
image_urls, video_urls = _extract_media_urls(message) image_urls, video_urls, file_urls, unknown_attachments = _extract_media_urls(message)
reply_mid, reply_preview = _extract_max_reply(message) reply_mid, reply_preview = _extract_max_reply(message)
return ParsedMessage( return ParsedMessage(
message_id=message_id, message_id=message_id,
@@ -125,6 +207,8 @@ def parse_message(message: Any) -> ParsedMessage:
text=text, text=text,
image_urls=image_urls, image_urls=image_urls,
video_urls=video_urls, video_urls=video_urls,
file_urls=file_urls,
unknown_attachments=unknown_attachments,
reply_to_max_message_id=reply_mid, reply_to_max_message_id=reply_mid,
reply_preview_text=reply_preview, reply_preview_text=reply_preview,
) )
+4
View File
@@ -10,6 +10,10 @@ class ParsedMessage:
text: str text: str
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)
# URL -> исходное имя файла (если удалось определить в MAX).
file_names_by_url: dict[str, str] = field(default_factory=dict)
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
reply_preview_text: str | None = None reply_preview_text: str | None = None
+88 -9
View File
@@ -64,6 +64,11 @@ def _is_supported_telegram_message(message: dict[str, Any]) -> bool:
if has_video: if has_video:
return True return True
for key in ("document", "audio", "voice", "animation", "sticker", "video_note"):
value = message.get(key)
if isinstance(value, dict) and value.get("file_id"):
return True
return False return False
@@ -199,11 +204,7 @@ class TelegramToMaxBridge:
chat_title = _telegram_chat_title(chat) chat_title = _telegram_chat_title(chat)
normalized = _normalize_title(chat_title) normalized = _normalize_title(chat_title)
if not normalized: max_chat_id = self._resolve_max_chat_id(message=message, normalized_title=normalized, chat=chat)
logger.error("Telegram chat without title/username, skip (chat=%s)", chat)
return
max_chat_id = self._resolve_max_chat_id_by_title(normalized)
if max_chat_id is None: if max_chat_id is None:
# требование: если в MAX нет канала/группы — ошибка и не пересылать # требование: если в MAX нет канала/группы — ошибка и не пересылать
logger.error("MAX чат с названием '%s' не найден — сообщение не пересылаю", chat_title) logger.error("MAX чат с названием '%s' не найден — сообщение не пересылаю", chat_title)
@@ -314,8 +315,13 @@ class TelegramToMaxBridge:
reply_to = self._resolve_reply_to_max_id(max_chat_id=max_chat_id, message=messages[0]) reply_to = self._resolve_reply_to_max_id(max_chat_id=max_chat_id, message=messages[0])
attachments: list[Any] = [] attachments: list[Any] = []
extra_links: list[str] = []
for m in messages: for m in messages:
attachments.extend(await self._extract_attachments(m)) extracted, links = await self._extract_attachments(m)
attachments.extend(extracted)
extra_links.extend(links)
text = self._append_file_links(text=text, links=extra_links)
if not text.strip() and not attachments: if not text.strip() and not attachments:
return return
@@ -375,7 +381,8 @@ class TelegramToMaxBridge:
) -> None: ) -> None:
raw_text = str(message.get("text") or message.get("caption") or "").strip() raw_text = str(message.get("text") or message.get("caption") or "").strip()
text = _format_forward_text(sender=message.get("from"), text=raw_text) text = _format_forward_text(sender=message.get("from"), text=raw_text)
attachments = await self._extract_attachments(message) attachments, extra_links = await self._extract_attachments(message)
text = self._append_file_links(text=text, links=extra_links)
if not text.strip() and not attachments: if not text.strip() and not attachments:
return return
@@ -430,8 +437,58 @@ class TelegramToMaxBridge:
# reply_to в MAX — это id сообщения; если не нашли, просто отправляем без reply # reply_to в MAX — это id сообщения; если не нашли, просто отправляем без reply
return mapped return mapped
async def _extract_attachments(self, message: dict[str, Any]) -> list[Any]: def _resolve_max_chat_id(
self,
*,
message: dict[str, Any],
normalized_title: str,
chat: dict[str, Any],
) -> int | None:
# При reply маршрутизируем по исходному сообщению (контекст диалога),
# чтобы не зависеть от username/title Telegram-чата.
from_reply = self._resolve_max_chat_id_from_reply(message)
if from_reply is not None:
return from_reply
if not normalized_title:
logger.error("Telegram chat without title/username, skip (chat=%s)", chat)
return None
return self._resolve_max_chat_id_by_title(normalized_title)
def _resolve_max_chat_id_from_reply(self, message: dict[str, Any]) -> int | None:
reply = message.get("reply_to_message")
if not isinstance(reply, dict):
return None
reply_mid = reply.get("message_id")
if reply_mid is None:
return None
chat = message.get("chat")
if not isinstance(chat, dict):
return None
telegram_chat_id = str(chat.get("id"))
max_chat_id = self._storage.get_max_chat_id_for_telegram(
telegram_chat_id=telegram_chat_id,
telegram_message_id=str(reply_mid),
)
if not max_chat_id:
return None
try:
return int(max_chat_id)
except ValueError:
logger.warning(
"Invalid max_chat_id '%s' in mapping for Telegram %s/%s",
max_chat_id,
telegram_chat_id,
reply_mid,
)
return None
async def _extract_attachments(self, message: dict[str, Any]) -> tuple[list[Any], list[str]]:
attachments: list[Any] = [] attachments: list[Any] = []
file_links: list[str] = []
# photo: массив размеров, берём последний (самый большой) # photo: массив размеров, берём последний (самый большой)
photos = message.get("photo") photos = message.get("photo")
@@ -459,7 +516,29 @@ class TelegramToMaxBridge:
except Exception: except Exception:
logger.exception("Cannot fetch Telegram video URL") logger.exception("Cannot fetch Telegram video URL")
return attachments for key in ("document", "audio", "voice", "animation", "sticker", "video_note"):
value = message.get(key)
if not isinstance(value, dict):
continue
file_id = str(value.get("file_id") or "")
if not file_id:
continue
try:
url = await self._telegram.get_file_url(file_id)
file_links.append(url)
except Exception:
logger.exception("Cannot fetch Telegram %s URL", key)
return attachments, file_links
@staticmethod
def _append_file_links(*, text: str, links: list[str]) -> str:
uniq_links = list(dict.fromkeys([str(link).strip() for link in links if str(link).strip()]))
if not uniq_links:
return text
links_block = "\n".join(f"- {url}" for url in uniq_links)
suffix = f"\n\n[Telegram files]\n{links_block}"
return f"{text}{suffix}" if text else suffix.strip()
def _refresh_max_chat_cache(self) -> None: def _refresh_max_chat_cache(self) -> None:
title_to_id: dict[str, int] = {} title_to_id: dict[str, int] = {}
+16
View File
@@ -105,6 +105,22 @@ class BridgeStorage:
return None return None
return str(row[0]) return str(row[0])
def get_max_chat_id_for_telegram(
self, *, telegram_chat_id: str, telegram_message_id: str
) -> str | None:
with closing(self._connect()) as conn:
row = conn.execute(
"""
SELECT max_chat_id
FROM message_mapping
WHERE telegram_chat_id = ? AND telegram_message_id = ?
""",
(telegram_chat_id, telegram_message_id),
).fetchone()
if not row:
return None
return str(row[0])
def get_telegram_message_id_for_max( def get_telegram_message_id_for_max(
self, *, telegram_chat_id: str, max_chat_id: str, max_message_id: str self, *, telegram_chat_id: str, max_chat_id: str, max_message_id: str
) -> str | None: ) -> str | None:
+169 -3
View File
@@ -1,11 +1,38 @@
import asyncio import asyncio
import json
import os
import pathlib
import tempfile
import urllib.parse
from typing import Any from typing import Any
import requests import requests
class TelegramApiError(RuntimeError): class TelegramApiError(RuntimeError):
pass def __init__(
self,
message: str,
*,
method: str | None = None,
status_code: int | None = None,
error_code: int | None = None,
description: str | None = None,
parameters: dict[str, Any] | None = None,
) -> None:
super().__init__(message)
self.method = method
self.status_code = status_code
self.error_code = error_code
self.description = description
self.parameters = parameters or {}
@property
def migrate_to_chat_id(self) -> str | None:
value = self.parameters.get("migrate_to_chat_id")
if value is None:
return None
return str(value)
class TelegramClient: class TelegramClient:
@@ -16,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:
@@ -80,6 +108,63 @@ class TelegramClient:
payload["reply_to_message_id"] = reply_to_message_id payload["reply_to_message_id"] = reply_to_message_id
return await self._request("sendVideo", payload) return await self._request("sendVideo", payload)
async def send_document(
self,
chat_id: str,
document_url: str,
file_name: str | None = None,
caption: str | None = None,
*,
reply_to_message_id: int | None = None,
) -> dict[str, Any]:
# Telegram часто не может скачать URL, которые доступны только клиенту MAX.
# Поэтому скачиваем сами во временный файл и отправляем как multipart upload.
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:
payload["caption"] = caption
if reply_to_message_id is not None:
payload["reply_to_message_id"] = str(int(reply_to_message_id))
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,
chat_id: str, chat_id: str,
@@ -208,11 +293,92 @@ 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
description: str | None = None
parameters: dict[str, Any] = {}
try:
payload_data = response.json()
if isinstance(payload_data, dict):
if isinstance(payload_data.get("error_code"), int):
error_code = payload_data.get("error_code")
if isinstance(payload_data.get("description"), str):
description = payload_data.get("description")
if isinstance(payload_data.get("parameters"), dict):
parameters = payload_data.get("parameters", {})
except (json.JSONDecodeError, ValueError):
payload_data = None
raise TelegramApiError( raise TelegramApiError(
f"Telegram HTTP error on {method}: {response.status_code} {response.text}" f"Telegram HTTP error on {method}: {response.status_code} {response.text}",
method=method,
status_code=response.status_code,
error_code=error_code,
description=description,
parameters=parameters,
) )
data = response.json() data = response.json()
if not data.get("ok"): if not data.get("ok"):
raise TelegramApiError(f"Telegram API error on {method}: {data}") raise TelegramApiError(
f"Telegram API error on {method}: {data}",
method=method,
error_code=data.get("error_code") if isinstance(data.get("error_code"), int) else None,
description=data.get("description") if isinstance(data.get("description"), str) 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
+33 -2
View File
@@ -90,6 +90,37 @@ def _normalize(value: str) -> str:
return str(value or "").strip().casefold() return str(value or "").strip().casefold()
def _deduplicate_chats(chats: list[Any]) -> list[Any]:
"""
Возвращает уникальные чаты с сохранением исходного порядка.
Сначала пытаемся уникализировать по chat.id, затем по нормализованному title.
"""
unique: list[Any] = []
seen_ids: set[str] = set()
seen_titles: set[str] = set()
for chat in chats:
chat_id = getattr(chat, "id", None)
if chat_id is not None:
key_id = str(chat_id).strip()
if key_id in seen_ids:
continue
seen_ids.add(key_id)
unique.append(chat)
continue
key_title = _normalize(_max_chat_title(chat))
if not key_title:
unique.append(chat)
continue
if key_title in seen_titles:
continue
seen_titles.add(key_title)
unique.append(chat)
return unique
async def _refresh_chats_best_effort(max_client: MaxClient) -> None: async def _refresh_chats_best_effort(max_client: MaxClient) -> None:
# group.py: fetch_chats(marker=None) заполняет max_client.chats # group.py: fetch_chats(marker=None) заполняет max_client.chats
try: try:
@@ -104,7 +135,7 @@ def _find_chat_by_title(max_client: MaxClient, title: str) -> Any | None:
wanted = _normalize(title) wanted = _normalize(title)
if not wanted: if not wanted:
return None return None
chats = list(getattr(max_client, "chats", []) or []) chats = _deduplicate_chats(list(getattr(max_client, "chats", []) or []))
for c in chats: for c in chats:
if _normalize(_max_chat_title(c)) == wanted: if _normalize(_max_chat_title(c)) == wanted:
return c return c
@@ -166,7 +197,7 @@ async def handle_control_command(
if cmd == "/list": if cmd == "/list":
await _refresh_chats_best_effort(max_client) await _refresh_chats_best_effort(max_client)
chats = list(getattr(max_client, "chats", []) or []) chats = _deduplicate_chats(list(getattr(max_client, "chats", []) or []))
if not chats: if not chats:
return "Список чатов пуст (или клиент MAX ещё не успел их загрузить)." return "Список чатов пуст (или клиент MAX ещё не успел их загрузить)."
+475
View File
@@ -0,0 +1,475 @@
# Техническое задание: двунаправленный мост 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;
- если MAX message является reply, не копировать медиа из `link.message` (исходного сообщения);
- если найден, отправлять `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 без повторов между рестартами).