Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d5ddde3103 | ||
|
|
cabfb9b5f9 | ||
|
|
cc9d05b11b | ||
|
|
64040d6bc5 | ||
|
|
87c4c765e5 | ||
|
|
2d32bf4807 | ||
|
|
d019e884b6 | ||
|
|
e40297787e | ||
|
|
dcf4441340 |
+307
-24
@@ -5,7 +5,7 @@ from typing import Any
|
||||
from max_parser import parse_message
|
||||
from models import ParsedMessage
|
||||
from pymax import MaxClient
|
||||
from pymax.types import PhotoAttach, VideoAttach
|
||||
from pymax.types import AudioAttach, FileAttach, Message, PhotoAttach, StickerAttach, VideoAttach
|
||||
from storage import BridgeStorage
|
||||
from telegram_api import TelegramApiError, TelegramClient
|
||||
|
||||
@@ -19,9 +19,9 @@ class MaxToTelegramBridge:
|
||||
self._storage = storage
|
||||
|
||||
async def forward_message(self, max_message: Any) -> None:
|
||||
if self._is_self_message(max_message):
|
||||
logger.debug("Skip self message %s/%s", getattr(max_message, "chat_id", "?"), getattr(max_message, "id", "?"))
|
||||
return
|
||||
#if self._is_self_message(max_message):
|
||||
# logger.debug("Skip self message %s/%s", getattr(max_message, "chat_id", "?"), getattr(max_message, "id", "?"))
|
||||
# return
|
||||
|
||||
parsed = parse_message(max_message)
|
||||
parsed = await self._enrich_from_max(max_message, parsed)
|
||||
@@ -34,11 +34,6 @@ class MaxToTelegramBridge:
|
||||
# - если найден целевой Telegram-чат (не 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()
|
||||
routed = self._storage.get_chat_route(max_chat_title_norm=normalized)
|
||||
if routed:
|
||||
@@ -64,6 +59,7 @@ class MaxToTelegramBridge:
|
||||
text=parsed.text,
|
||||
include_chat_name=is_fallback,
|
||||
)
|
||||
text = self._append_unknown_attachment_notice(parsed=parsed, text=text)
|
||||
|
||||
reply_telegram_mid = self._resolve_telegram_reply_to(
|
||||
telegram_chat_id=str(target_chat_id),
|
||||
@@ -73,12 +69,12 @@ class MaxToTelegramBridge:
|
||||
if parsed.reply_to_max_message_id and reply_telegram_mid is None:
|
||||
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:
|
||||
logger.debug("Skip empty message %s/%s", parsed.chat_id, parsed.message_id)
|
||||
return
|
||||
text = self._build_fallback_unknown_notice(parsed)
|
||||
|
||||
total_media = len(parsed.image_urls) + len(parsed.video_urls)
|
||||
sent_any = False
|
||||
if total_media > 1:
|
||||
# Отправляем одним альбомом в Telegram (единое сообщение).
|
||||
target_chat_id, sent_messages = await self._send_with_migration_retry(
|
||||
@@ -103,6 +99,7 @@ class MaxToTelegramBridge:
|
||||
max_chat_id=str(parsed.chat_id),
|
||||
max_message_id=str(parsed.message_id),
|
||||
)
|
||||
sent_any = True
|
||||
self._storage.mark_forwarded(parsed.message_id, parsed.chat_id)
|
||||
logger.info(
|
||||
"Forwarded media group %s/%s (images=%s, videos=%s)",
|
||||
@@ -113,8 +110,8 @@ class MaxToTelegramBridge:
|
||||
)
|
||||
return
|
||||
|
||||
sent_any = False
|
||||
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=target_chat_id,
|
||||
max_chat_title_norm=normalized,
|
||||
@@ -179,15 +176,83 @@ class MaxToTelegramBridge:
|
||||
)
|
||||
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
|
||||
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
|
||||
|
||||
self._storage.mark_forwarded(parsed.message_id, parsed.chat_id)
|
||||
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.message_id,
|
||||
len(parsed.image_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,
|
||||
*,
|
||||
@@ -277,27 +342,172 @@ class MaxToTelegramBridge:
|
||||
except Exception:
|
||||
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:
|
||||
if isinstance(attach, PhotoAttach):
|
||||
parsed.image_urls.extend(self._extract_photo_urls(attach))
|
||||
elif isinstance(attach, VideoAttach):
|
||||
continue
|
||||
|
||||
if isinstance(attach, VideoAttach):
|
||||
try:
|
||||
video = await self._max_client.get_video_by_id(
|
||||
chat_id=max_message.chat_id,
|
||||
message_id=max_message.id,
|
||||
chat_id=message_chat_id,
|
||||
message_id=message_id,
|
||||
video_id=attach.video_id,
|
||||
)
|
||||
video_url = getattr(video, "url", None)
|
||||
if video_url:
|
||||
parsed.video_urls.append(str(video_url))
|
||||
except Exception:
|
||||
logger.exception("Cannot resolve video URL from Max")
|
||||
logger.exception("Cannot resolve video URL from Max (%s)", source_tag)
|
||||
continue
|
||||
|
||||
# Убираем дубли URL, если парсер и enrich нашли одинаковые вложения.
|
||||
parsed.image_urls = list(dict.fromkeys(parsed.image_urls))
|
||||
parsed.video_urls = list(dict.fromkeys(parsed.video_urls))
|
||||
return parsed
|
||||
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:
|
||||
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:
|
||||
sender = getattr(max_message, "sender", None)
|
||||
@@ -345,3 +555,76 @@ class MaxToTelegramBridge:
|
||||
|
||||
walk(attach)
|
||||
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
@@ -72,8 +72,9 @@ def main() -> None:
|
||||
health.mark_max_event()
|
||||
try:
|
||||
await bridge.forward_message(message)
|
||||
except Exception:
|
||||
except Exception as exc:
|
||||
logger.exception("Failed to forward Max message")
|
||||
await bridge.notify_delivery_failure(message, exc)
|
||||
|
||||
asyncio.run(max_client.start())
|
||||
|
||||
|
||||
+87
-3
@@ -36,15 +36,76 @@ def _is_video(media_type: str) -> bool:
|
||||
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] = []
|
||||
video_urls: list[str] = []
|
||||
file_urls: list[str] = []
|
||||
unknown_attachments: list[str] = []
|
||||
|
||||
# В PyMax рабочее поле для вложений обычно называется attaches.
|
||||
raw_attachments = _get_attr(message, ["attaches", "attachments", "media", "files"], default=[]) or []
|
||||
for item in raw_attachments:
|
||||
data = _as_dict(item)
|
||||
media_type = _stringify(data.get("type") or data.get("media_type") or data.get("kind"))
|
||||
is_forward_like = _is_forward_like(data)
|
||||
url = _stringify(
|
||||
data.get("base_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"))
|
||||
|
||||
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
|
||||
|
||||
if _is_image(media_type):
|
||||
image_urls.append(url)
|
||||
elif _is_video(media_type):
|
||||
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]:
|
||||
@@ -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"
|
||||
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)
|
||||
return ParsedMessage(
|
||||
message_id=message_id,
|
||||
@@ -125,6 +207,8 @@ def parse_message(message: Any) -> ParsedMessage:
|
||||
text=text,
|
||||
image_urls=image_urls,
|
||||
video_urls=video_urls,
|
||||
file_urls=file_urls,
|
||||
unknown_attachments=unknown_attachments,
|
||||
reply_to_max_message_id=reply_mid,
|
||||
reply_preview_text=reply_preview,
|
||||
)
|
||||
|
||||
@@ -10,6 +10,10 @@ class ParsedMessage:
|
||||
text: str
|
||||
image_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 указывает на исходное сообщение (тред).
|
||||
reply_to_max_message_id: str | None = None
|
||||
reply_preview_text: str | None = None
|
||||
|
||||
+38
-4
@@ -64,6 +64,11 @@ def _is_supported_telegram_message(message: dict[str, Any]) -> bool:
|
||||
if has_video:
|
||||
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
|
||||
|
||||
|
||||
@@ -314,8 +319,13 @@ class TelegramToMaxBridge:
|
||||
reply_to = self._resolve_reply_to_max_id(max_chat_id=max_chat_id, message=messages[0])
|
||||
|
||||
attachments: list[Any] = []
|
||||
extra_links: list[str] = []
|
||||
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:
|
||||
return
|
||||
@@ -375,7 +385,8 @@ class TelegramToMaxBridge:
|
||||
) -> None:
|
||||
raw_text = str(message.get("text") or message.get("caption") or "").strip()
|
||||
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:
|
||||
return
|
||||
|
||||
@@ -430,8 +441,9 @@ class TelegramToMaxBridge:
|
||||
# reply_to в MAX — это id сообщения; если не нашли, просто отправляем без reply
|
||||
return mapped
|
||||
|
||||
async def _extract_attachments(self, message: dict[str, Any]) -> list[Any]:
|
||||
async def _extract_attachments(self, message: dict[str, Any]) -> tuple[list[Any], list[str]]:
|
||||
attachments: list[Any] = []
|
||||
file_links: list[str] = []
|
||||
|
||||
# photo: массив размеров, берём последний (самый большой)
|
||||
photos = message.get("photo")
|
||||
@@ -459,7 +471,29 @@ class TelegramToMaxBridge:
|
||||
except Exception:
|
||||
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:
|
||||
title_to_id: dict[str, int] = {}
|
||||
|
||||
@@ -1,5 +1,9 @@
|
||||
import asyncio
|
||||
import json
|
||||
import os
|
||||
import pathlib
|
||||
import tempfile
|
||||
import urllib.parse
|
||||
from typing import Any
|
||||
|
||||
import requests
|
||||
@@ -39,6 +43,7 @@ class TelegramClient:
|
||||
self._timeout = timeout
|
||||
self._chat_title_to_id: dict[str, str] = {}
|
||||
self._me: dict[str, Any] | None = None
|
||||
self._tmp_root: str | None = None
|
||||
|
||||
@property
|
||||
def fallback_user_id(self) -> str:
|
||||
@@ -103,6 +108,63 @@ class TelegramClient:
|
||||
payload["reply_to_message_id"] = reply_to_message_id
|
||||
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(
|
||||
self,
|
||||
chat_id: str,
|
||||
@@ -231,6 +293,9 @@ class TelegramClient:
|
||||
return requests.post(url, json=payload, timeout=self._timeout)
|
||||
|
||||
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:
|
||||
error_code: int | None = None
|
||||
description: str | None = None
|
||||
@@ -265,3 +330,55 @@ class TelegramClient:
|
||||
parameters=data.get("parameters") if isinstance(data.get("parameters"), dict) else None,
|
||||
)
|
||||
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
@@ -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 без повторов между рестартами).
|
||||
|
||||
Reference in New Issue
Block a user