Compare commits
12
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2d32bf4807 | ||
|
|
d019e884b6 | ||
|
|
e40297787e | ||
|
|
dcf4441340 | ||
|
|
1ef5cf4e38 | ||
|
|
1fd5e4ea39 | ||
|
|
bcf9366cf6 | ||
|
|
102009369a | ||
|
|
7a3c7317ea | ||
|
|
7d522e5f73 | ||
|
|
70f1cea78e | ||
|
|
763ca45ab2 |
@@ -0,0 +1,78 @@
|
|||||||
|
# Сборка образа и публикация на Docker Hub (hub.docker.com).
|
||||||
|
#
|
||||||
|
# Настройка в GitHub → Settings → Secrets and variables → Actions:
|
||||||
|
# DOCKERHUB_USERNAME — логин Docker Hub
|
||||||
|
# DOCKERHUB_TOKEN — Access Token (рекомендуется), см. https://hub.docker.com/settings/security
|
||||||
|
#
|
||||||
|
# Опционально: Variables → DOCKERHUB_IMAGE (например org/max2telegram), если имя образа
|
||||||
|
# отличается от <DOCKERHUB_USERNAME>/max2telegram.
|
||||||
|
#
|
||||||
|
# Теги образа:
|
||||||
|
# push в master → dev (+ sha)
|
||||||
|
# push тега v* → {{version}}, {{major}}.{{minor}}, latest (+ sha)
|
||||||
|
|
||||||
|
name: Docker Hub
|
||||||
|
|
||||||
|
on:
|
||||||
|
push:
|
||||||
|
branches:
|
||||||
|
- main
|
||||||
|
- master
|
||||||
|
tags:
|
||||||
|
- 'v*'
|
||||||
|
workflow_dispatch:
|
||||||
|
|
||||||
|
concurrency:
|
||||||
|
group: docker-hub-${{ github.ref }}
|
||||||
|
cancel-in-progress: true
|
||||||
|
|
||||||
|
permissions:
|
||||||
|
contents: read
|
||||||
|
actions: write # кэш слоёв Buildx (type=gha)
|
||||||
|
|
||||||
|
jobs:
|
||||||
|
build-and-push:
|
||||||
|
runs-on: ubuntu-latest
|
||||||
|
steps:
|
||||||
|
- name: Checkout
|
||||||
|
uses: actions/checkout@v4
|
||||||
|
|
||||||
|
- name: Set image name
|
||||||
|
id: image
|
||||||
|
run: |
|
||||||
|
if [ -n "${{ vars.DOCKERHUB_IMAGE }}" ]; then
|
||||||
|
echo "name=${{ vars.DOCKERHUB_IMAGE }}" >> "$GITHUB_OUTPUT"
|
||||||
|
else
|
||||||
|
echo "name=${{ secrets.DOCKERHUB_USERNAME }}/max2telegram" >> "$GITHUB_OUTPUT"
|
||||||
|
fi
|
||||||
|
|
||||||
|
- name: Docker Hub login
|
||||||
|
uses: docker/login-action@v3
|
||||||
|
with:
|
||||||
|
username: ${{ secrets.DOCKERHUB_USERNAME }}
|
||||||
|
password: ${{ secrets.DOCKERHUB_TOKEN }}
|
||||||
|
|
||||||
|
- name: Set up Buildx
|
||||||
|
uses: docker/setup-buildx-action@v3
|
||||||
|
|
||||||
|
- name: Docker metadata (теги)
|
||||||
|
id: meta
|
||||||
|
uses: docker/metadata-action@v5
|
||||||
|
with:
|
||||||
|
images: ${{ steps.image.outputs.name }}
|
||||||
|
tags: |
|
||||||
|
type=raw,value=dev,enable=${{ github.ref == 'refs/heads/master' }}
|
||||||
|
type=semver,pattern={{version}}
|
||||||
|
type=semver,pattern={{major}}.{{minor}}
|
||||||
|
type=raw,value=latest,enable=${{ startsWith(github.ref, 'refs/tags/') }}
|
||||||
|
|
||||||
|
- name: Build and push
|
||||||
|
uses: docker/build-push-action@v6
|
||||||
|
with:
|
||||||
|
context: .
|
||||||
|
file: ./Dockerfile
|
||||||
|
push: true
|
||||||
|
tags: ${{ steps.meta.outputs.tags }}
|
||||||
|
labels: ${{ steps.meta.outputs.labels }}
|
||||||
|
cache-from: type=gha
|
||||||
|
cache-to: type=gha,mode=max
|
||||||
+3
-1
@@ -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
|
||||||
|
|
||||||
|
|||||||
+339
-33
@@ -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,60 @@ 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
|
||||||
|
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,
|
||||||
|
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 +222,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 +340,152 @@ 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:
|
||||||
|
if not (parsed.text or "").strip():
|
||||||
|
linked_text = str(getattr(linked_message, "text", "") or "").strip()
|
||||||
|
if linked_text:
|
||||||
|
parsed.text = linked_text
|
||||||
|
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.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)
|
||||||
|
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:
|
||||||
|
parsed.file_urls.append(str(url))
|
||||||
|
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 +533,71 @@ 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})."
|
||||||
|
|||||||
+2
-1
@@ -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
@@ -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,
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -10,6 +10,8 @@ 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)
|
||||||
|
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
|
||||||
|
|||||||
+80
-13
@@ -48,6 +48,30 @@ def _format_forward_text(*, sender: dict[str, Any] | None, text: str) -> str:
|
|||||||
return header
|
return header
|
||||||
|
|
||||||
|
|
||||||
|
def _is_supported_telegram_message(message: dict[str, Any]) -> bool:
|
||||||
|
# Текстовые сообщения и команды.
|
||||||
|
text = message.get("text")
|
||||||
|
if isinstance(text, str) and text.strip():
|
||||||
|
return True
|
||||||
|
|
||||||
|
photos = message.get("photo")
|
||||||
|
has_photo = isinstance(photos, list) and any(isinstance(p, dict) and p.get("file_id") for p in photos)
|
||||||
|
if has_photo:
|
||||||
|
return True
|
||||||
|
|
||||||
|
video = message.get("video")
|
||||||
|
has_video = isinstance(video, dict) and video.get("file_id")
|
||||||
|
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
|
||||||
|
|
||||||
|
|
||||||
@dataclass
|
@dataclass
|
||||||
class _MediaGroupBuffer:
|
class _MediaGroupBuffer:
|
||||||
first_seen_monotonic: float
|
first_seen_monotonic: float
|
||||||
@@ -118,7 +142,7 @@ class TelegramToMaxBridge:
|
|||||||
max_update_id = upd_id if max_update_id is None else max(max_update_id, upd_id)
|
max_update_id = upd_id if max_update_id is None else max(max_update_id, upd_id)
|
||||||
|
|
||||||
message = None
|
message = None
|
||||||
for container in ("message", "edited_message", "channel_post", "edited_channel_post"):
|
for container in ("message", "channel_post"):
|
||||||
candidate = upd.get(container)
|
candidate = upd.get(container)
|
||||||
if isinstance(candidate, dict):
|
if isinstance(candidate, dict):
|
||||||
message = candidate
|
message = candidate
|
||||||
@@ -126,6 +150,9 @@ class TelegramToMaxBridge:
|
|||||||
if not message:
|
if not message:
|
||||||
continue
|
continue
|
||||||
|
|
||||||
|
if not _is_supported_telegram_message(message):
|
||||||
|
continue
|
||||||
|
|
||||||
if self._is_own_telegram_message(message):
|
if self._is_own_telegram_message(message):
|
||||||
continue
|
continue
|
||||||
|
|
||||||
@@ -151,6 +178,14 @@ class TelegramToMaxBridge:
|
|||||||
if not isinstance(chat, dict):
|
if not isinstance(chat, dict):
|
||||||
return
|
return
|
||||||
|
|
||||||
|
# Команда привязки чата Telegram к названию чата в MAX (для Max->Telegram маршрутизации).
|
||||||
|
# Обрабатываем раньше control-команд, чтобы /bind_max не попадала как "неизвестная".
|
||||||
|
text = str(message.get("text") or "").strip()
|
||||||
|
cmd = text.split(maxsplit=1)[0].split("@", 1)[0].strip().casefold() if text else ""
|
||||||
|
if cmd == "/bind_max":
|
||||||
|
await self._handle_bind_max_command(message, chat)
|
||||||
|
return
|
||||||
|
|
||||||
# Управление MAX через Telegram: только личка боту и только от fallback_user_id.
|
# Управление MAX через Telegram: только личка боту и только от fallback_user_id.
|
||||||
# В этом случае команду не пересылаем в MAX.
|
# В этом случае команду не пересылаем в MAX.
|
||||||
try:
|
try:
|
||||||
@@ -167,13 +202,6 @@ class TelegramToMaxBridge:
|
|||||||
logger.exception("Cannot send Telegram reply for control command")
|
logger.exception("Cannot send Telegram reply for control command")
|
||||||
return
|
return
|
||||||
|
|
||||||
# Команда привязки чата Telegram к названию чата в MAX (для Max->Telegram маршрутизации).
|
|
||||||
# Работает даже при privacy mode, т.к. команды приходят боту.
|
|
||||||
text = str(message.get("text") or "").strip()
|
|
||||||
if text.startswith("/bind_max"):
|
|
||||||
await self._handle_bind_max_command(message, chat)
|
|
||||||
return
|
|
||||||
|
|
||||||
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:
|
if not normalized:
|
||||||
@@ -209,10 +237,13 @@ class TelegramToMaxBridge:
|
|||||||
|
|
||||||
async def _handle_bind_max_command(self, message: dict[str, Any], chat: dict[str, Any]) -> None:
|
async def _handle_bind_max_command(self, message: dict[str, Any], chat: dict[str, Any]) -> None:
|
||||||
raw = str(message.get("text") or "").strip()
|
raw = str(message.get("text") or "").strip()
|
||||||
|
chat_id = str(chat.get("id") or "")
|
||||||
# формат: /bind_max <точное название чата в MAX>
|
# формат: /bind_max <точное название чата в MAX>
|
||||||
parts = raw.split(maxsplit=1)
|
parts = raw.split(maxsplit=1)
|
||||||
if len(parts) < 2 or not parts[1].strip():
|
if len(parts) < 2 or not parts[1].strip():
|
||||||
logger.error("bind_max: missing MAX chat title. Use: /bind_max <MAX chat title>")
|
logger.error("bind_max: missing MAX chat title. Use: /bind_max <MAX chat title>")
|
||||||
|
if chat_id:
|
||||||
|
await self._telegram.send_text(chat_id=chat_id, text="Использование: /bind_max <точное название чата в MAX>")
|
||||||
return
|
return
|
||||||
|
|
||||||
max_title = parts[1].strip()
|
max_title = parts[1].strip()
|
||||||
@@ -222,9 +253,11 @@ class TelegramToMaxBridge:
|
|||||||
max_chat_id = self._resolve_max_chat_id_by_title(norm)
|
max_chat_id = self._resolve_max_chat_id_by_title(norm)
|
||||||
if max_chat_id is None:
|
if max_chat_id is None:
|
||||||
logger.error("bind_max: MAX чат '%s' не найден — привязку не сохраняю", max_title)
|
logger.error("bind_max: MAX чат '%s' не найден — привязку не сохраняю", max_title)
|
||||||
|
if chat_id:
|
||||||
|
await self._telegram.send_text(chat_id=chat_id, text=f"MAX чат '{max_title}' не найден. Привязка не сохранена.")
|
||||||
return
|
return
|
||||||
|
|
||||||
telegram_chat_id = str(chat.get("id"))
|
telegram_chat_id = chat_id
|
||||||
telegram_title = _telegram_chat_title(chat)
|
telegram_title = _telegram_chat_title(chat)
|
||||||
self._storage.set_chat_route(
|
self._storage.set_chat_route(
|
||||||
max_chat_title_norm=norm,
|
max_chat_title_norm=norm,
|
||||||
@@ -237,6 +270,11 @@ class TelegramToMaxBridge:
|
|||||||
telegram_title,
|
telegram_title,
|
||||||
telegram_chat_id,
|
telegram_chat_id,
|
||||||
)
|
)
|
||||||
|
if chat_id:
|
||||||
|
await self._telegram.send_text(
|
||||||
|
chat_id=chat_id,
|
||||||
|
text=f"Канал успешно привязан: Telegram '{telegram_title or telegram_chat_id}' -> MAX '{max_title}'.",
|
||||||
|
)
|
||||||
|
|
||||||
async def _flush_ready_media_groups(self) -> None:
|
async def _flush_ready_media_groups(self) -> None:
|
||||||
now = time.monotonic()
|
now = time.monotonic()
|
||||||
@@ -281,8 +319,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
|
||||||
@@ -342,7 +385,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
|
||||||
|
|
||||||
@@ -397,8 +441,9 @@ 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]:
|
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")
|
||||||
@@ -426,7 +471,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] = {}
|
||||||
|
|||||||
+73
-6
@@ -1,11 +1,34 @@
|
|||||||
import asyncio
|
import asyncio
|
||||||
|
import json
|
||||||
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:
|
||||||
@@ -80,6 +103,24 @@ 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,
|
||||||
|
caption: str | None = None,
|
||||||
|
*,
|
||||||
|
reply_to_message_id: int | None = None,
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
payload: dict[str, Any] = {
|
||||||
|
"chat_id": chat_id,
|
||||||
|
"document": document_url,
|
||||||
|
}
|
||||||
|
if caption:
|
||||||
|
payload["caption"] = caption
|
||||||
|
if reply_to_message_id is not None:
|
||||||
|
payload["reply_to_message_id"] = reply_to_message_id
|
||||||
|
return await self._request("sendDocument", payload)
|
||||||
|
|
||||||
async def send_media_group(
|
async def send_media_group(
|
||||||
self,
|
self,
|
||||||
chat_id: str,
|
chat_id: str,
|
||||||
@@ -132,8 +173,8 @@ class TelegramClient:
|
|||||||
payload: dict[str, Any] = {
|
payload: dict[str, Any] = {
|
||||||
"timeout": timeout,
|
"timeout": timeout,
|
||||||
"limit": limit,
|
"limit": limit,
|
||||||
# чтобы получать посты из каналов (channel_post) и обычные сообщения
|
# Реагируем только на новые сообщения/посты (текст, фото, видео).
|
||||||
"allowed_updates": ["message", "edited_message", "channel_post", "edited_channel_post"],
|
"allowed_updates": ["message", "channel_post"],
|
||||||
}
|
}
|
||||||
if offset is not None:
|
if offset is not None:
|
||||||
payload["offset"] = offset
|
payload["offset"] = offset
|
||||||
@@ -145,7 +186,7 @@ class TelegramClient:
|
|||||||
# Важно: не делаем getUpdates нигде больше (иначе 409 Conflict).
|
# Важно: не делаем getUpdates нигде больше (иначе 409 Conflict).
|
||||||
# Наполняем кэш чатов только из этого потока.
|
# Наполняем кэш чатов только из этого потока.
|
||||||
for upd in updates:
|
for upd in updates:
|
||||||
for container in ("message", "edited_message", "channel_post", "edited_channel_post"):
|
for container in ("message", "channel_post"):
|
||||||
msg = upd.get(container)
|
msg = upd.get(container)
|
||||||
if not isinstance(msg, dict):
|
if not isinstance(msg, dict):
|
||||||
continue
|
continue
|
||||||
@@ -209,10 +250,36 @@ class TelegramClient:
|
|||||||
|
|
||||||
response = await asyncio.to_thread(_do_request)
|
response = await asyncio.to_thread(_do_request)
|
||||||
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
|
||||||
|
|||||||
@@ -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
|
||||||
@@ -159,13 +190,14 @@ async def handle_control_command(
|
|||||||
"/list — список активных чатов MAX\n"
|
"/list — список активных чатов MAX\n"
|
||||||
"/join <LINK> — присоединиться к группе/каналу по ссылке\n"
|
"/join <LINK> — присоединиться к группе/каналу по ссылке\n"
|
||||||
"/leave <НАЗВАНИЕ> — покинуть указанный канал\n"
|
"/leave <НАЗВАНИЕ> — покинуть указанный канал\n"
|
||||||
"/last_messages <НАЗВАНИЕ> — последние 10 сообщений из канала"
|
"/last_messages <НАЗВАНИЕ> — последние 10 сообщений из канала\n"
|
||||||
|
"/bind_max <НАЗВАНИЕ> — привязать текущий Telegram-чат к чату MAX"
|
||||||
)
|
)
|
||||||
|
|
||||||
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 ещё не успел их загрузить)."
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user