Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
40e85abf5a | ||
|
|
3a50472ca9 | ||
|
|
2b5260049b | ||
|
|
aecfe4dedf |
+1
-7
@@ -3,10 +3,4 @@ MAX_DEVICE_ID=a1b2c3d4-e5f6-7890-abcd-ef1234567890
|
||||
TG_BOT_TOKEN=123456:ABC-DEF1234...
|
||||
TG_FORUM_CHANNEL_ID=-1001234567890
|
||||
FALLBACK_USER_ID=987654321
|
||||
DATABASE_URL=sqlite+aiosqlite:///app/data/bridge.db
|
||||
REDIS_URL=redis://redis:6379/0
|
||||
TG_RATE_LIMIT_DELAY_SEC=3.5
|
||||
MAX_RATE_LIMIT_DELAY_SEC=1.0
|
||||
LS_TOPIC_PREFIX=👤
|
||||
MAX_RECONNECT_FETCH_LIMIT=50
|
||||
LOG_LEVEL=INFO
|
||||
DATABASE_URL=sqlite+aiosqlite:///data/bridge.db
|
||||
|
||||
+7
-4
@@ -2,6 +2,7 @@ from pathlib import Path
|
||||
|
||||
from pydantic import model_validator
|
||||
from pydantic_settings import BaseSettings, SettingsConfigDict
|
||||
from sqlalchemy.engine.url import make_url
|
||||
|
||||
|
||||
class Settings(BaseSettings):
|
||||
@@ -12,7 +13,7 @@ class Settings(BaseSettings):
|
||||
tg_bot_token: str
|
||||
tg_forum_channel_id: int
|
||||
fallback_user_id: int
|
||||
database_url: str = "sqlite+aiosqlite:///app/data/bridge.db"
|
||||
database_url: str = "sqlite+aiosqlite:///data/bridge.db"
|
||||
redis_url: str = "redis://redis:6379/0"
|
||||
tg_rate_limit_delay_sec: float = 3.5
|
||||
max_rate_limit_delay_sec: float = 1.0
|
||||
@@ -24,9 +25,11 @@ class Settings(BaseSettings):
|
||||
|
||||
@property
|
||||
def sqlite_path(self) -> str:
|
||||
if self.database_url.startswith("sqlite+aiosqlite:///"):
|
||||
return self.database_url.removeprefix("sqlite+aiosqlite:///")
|
||||
return "app/data/bridge.db"
|
||||
db = make_url(self.database_url).database
|
||||
if not db:
|
||||
return str(Path("data/bridge.db").resolve())
|
||||
path = Path(db)
|
||||
return str(path if path.is_absolute() else path.resolve())
|
||||
|
||||
@model_validator(mode="after")
|
||||
def _default_data_dir(self) -> "Settings":
|
||||
|
||||
+26
@@ -50,11 +50,37 @@ async def main() -> None:
|
||||
storage = SqliteStorage(session_factory)
|
||||
await storage.init()
|
||||
mappings = await storage.list_mappings()
|
||||
sqlite_file = Path(settings.sqlite_path)
|
||||
db_exists = sqlite_file.is_file()
|
||||
db_size_bytes = sqlite_file.stat().st_size if db_exists else 0
|
||||
logger.info(
|
||||
"storage_ready",
|
||||
sqlite_path=settings.sqlite_path,
|
||||
database_url=settings.database_url,
|
||||
db_file_exists=db_exists,
|
||||
db_file_size_bytes=db_size_bytes,
|
||||
mapping_count=len(mappings),
|
||||
)
|
||||
if len(mappings) == 0:
|
||||
logger.info(
|
||||
"storage_no_mappings",
|
||||
reason="chat_mappings_table_empty",
|
||||
effect="new_tg_forum_topic_will_be_created_for_each_max_chat",
|
||||
hint="after_container_recreate_check_volume_mount_data_dir_and_sqlite_path",
|
||||
sqlite_path=settings.sqlite_path,
|
||||
)
|
||||
|
||||
legacy_db = Path.cwd() / "app" / "data" / "bridge.db"
|
||||
if (
|
||||
legacy_db.is_file()
|
||||
and legacy_db.resolve() != sqlite_file.resolve()
|
||||
):
|
||||
logger.warning(
|
||||
"sqlite_legacy_path_detected",
|
||||
legacy_path=str(legacy_db.resolve()),
|
||||
active_path=settings.sqlite_path,
|
||||
hint="old_database_url_app/data/bridge_db_writes_outside_volume_use_data/bridge_db",
|
||||
)
|
||||
|
||||
queue = RedisQueue(settings.redis_url)
|
||||
await queue.connect()
|
||||
|
||||
@@ -3,7 +3,7 @@ from typing import Any
|
||||
import structlog
|
||||
from pymax.types.domain.message import Message as MaxMessage
|
||||
|
||||
from app.models.domain import MaxIncomingMessage
|
||||
from app.models.domain import MaxChatMeta, MaxIncomingMessage
|
||||
|
||||
logger = structlog.get_logger(__name__)
|
||||
|
||||
@@ -16,6 +16,21 @@ def get_forward_link(message: MaxMessage) -> dict[str, Any] | None:
|
||||
return link
|
||||
|
||||
|
||||
def get_reply_target(message: MaxMessage) -> int | None:
|
||||
link = getattr(message, "link", None)
|
||||
if not isinstance(link, dict):
|
||||
return None
|
||||
if str(link.get("type", "")).upper() != "REPLY":
|
||||
return None
|
||||
target = link.get("messageId") or link.get("message_id")
|
||||
if target is None:
|
||||
return None
|
||||
try:
|
||||
return int(target)
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def extract_forwarded_content(
|
||||
message: MaxMessage,
|
||||
) -> tuple[str, list[dict[str, Any]]] | None:
|
||||
@@ -196,14 +211,72 @@ async def resolve_media(client, message: MaxMessage) -> list[dict]:
|
||||
return items
|
||||
|
||||
|
||||
def build_chat_title(chat, ls_prefix: str) -> tuple[str, bool]:
|
||||
def _pick_image_url(*values: str | None) -> str | None:
|
||||
for value in values:
|
||||
if value and value.strip():
|
||||
return value.strip()
|
||||
return None
|
||||
|
||||
|
||||
def _chat_icon_url(chat) -> str | None:
|
||||
return _pick_image_url(
|
||||
getattr(chat, "base_icon_url", None),
|
||||
getattr(chat, "base_raw_icon_url", None),
|
||||
)
|
||||
|
||||
|
||||
def _user_icon_url(user) -> str | None:
|
||||
return _pick_image_url(
|
||||
getattr(user, "base_url", None),
|
||||
getattr(user, "base_raw_url", None),
|
||||
)
|
||||
|
||||
|
||||
async def resolve_chat_icon_url(
|
||||
client,
|
||||
chat,
|
||||
my_user_id: int | None,
|
||||
) -> str | None:
|
||||
is_dm = bool(getattr(chat, "is_dialog", False) or chat.type == "DIALOG")
|
||||
if is_dm:
|
||||
title = chat.title or "Контакт"
|
||||
if not title.startswith(ls_prefix.strip()):
|
||||
title = f"{ls_prefix}{title}"
|
||||
return title, True
|
||||
return chat.title or f"Чат {chat.id}", False
|
||||
for user_id in chat.participants or {}:
|
||||
if my_user_id and user_id == my_user_id:
|
||||
continue
|
||||
try:
|
||||
user = await client.get_user(user_id)
|
||||
icon_url = _user_icon_url(user)
|
||||
if icon_url:
|
||||
return icon_url
|
||||
except Exception:
|
||||
logger.debug(
|
||||
"max_dm_icon_lookup_failed",
|
||||
chat_id=chat.id,
|
||||
user_id=user_id,
|
||||
)
|
||||
return None
|
||||
|
||||
return _chat_icon_url(chat)
|
||||
|
||||
|
||||
def extract_chat_meta(chat, ls_prefix: str) -> MaxChatMeta:
|
||||
is_dm = bool(getattr(chat, "is_dialog", False) or chat.type == "DIALOG")
|
||||
chat_name = chat.title or ("Контакт" if is_dm else f"Чат {chat.id}")
|
||||
topic_title = chat_name
|
||||
if is_dm and not topic_title.startswith(ls_prefix.strip()):
|
||||
topic_title = f"{ls_prefix}{topic_title}"
|
||||
return MaxChatMeta(
|
||||
topic_title=topic_title,
|
||||
chat_name=chat_name,
|
||||
is_dm=is_dm,
|
||||
icon_url=_chat_icon_url(chat),
|
||||
participants_count=getattr(chat, "participants_count", 0) or 0,
|
||||
link=getattr(chat, "link", None),
|
||||
)
|
||||
|
||||
|
||||
def build_chat_title(chat, ls_prefix: str) -> tuple[str, bool]:
|
||||
meta = extract_chat_meta(chat, ls_prefix)
|
||||
return meta.topic_title, meta.is_dm
|
||||
|
||||
|
||||
def resolve_sender_name(user) -> str:
|
||||
|
||||
+40
-18
@@ -1,17 +1,20 @@
|
||||
import asyncio
|
||||
from pathlib import Path
|
||||
from uuid import uuid4
|
||||
|
||||
import aiohttp
|
||||
|
||||
import structlog
|
||||
from pymax import ExtraConfig, Message, WebClient
|
||||
from pymax.types.domain.enums import ChatType
|
||||
|
||||
from app.config import Settings
|
||||
from app.media_transfer import download_max_media, tmp_dir
|
||||
from app.media_transfer import download_max_image_url, download_max_media, tmp_dir
|
||||
from app.max_layer.client_holder import MaxClientHolder
|
||||
from app.max_layer.formatter import (
|
||||
build_chat_title,
|
||||
extract_chat_meta,
|
||||
extract_forwarded_content,
|
||||
resolve_chat_icon_url,
|
||||
format_forwarded_text,
|
||||
get_reply_target,
|
||||
resolve_media,
|
||||
resolve_raw_attaches,
|
||||
resolve_sender_name,
|
||||
@@ -85,6 +88,21 @@ class MaxListener:
|
||||
message_id=message.id,
|
||||
)
|
||||
|
||||
async def _download_chat_icon(self, icon_url: str | None) -> str | None:
|
||||
if not icon_url:
|
||||
return None
|
||||
dest = self._tmp_dir / f"chat_icon_{uuid4().hex}.jpg"
|
||||
try:
|
||||
async with aiohttp.ClientSession() as session:
|
||||
if not await download_max_image_url(session, icon_url, dest):
|
||||
logger.warning("max_chat_icon_download_failed", url=icon_url)
|
||||
return None
|
||||
return str(dest)
|
||||
except Exception:
|
||||
logger.exception("max_chat_icon_download_failed", url=icon_url)
|
||||
dest.unlink(missing_ok=True)
|
||||
return None
|
||||
|
||||
async def run(self) -> None:
|
||||
self._client = self.build_client()
|
||||
await self._client.start()
|
||||
@@ -103,6 +121,8 @@ class MaxListener:
|
||||
continue
|
||||
logger.info("max_history_fetched", chat_id=chat.id, count=len(messages))
|
||||
for msg in sorted(messages, key=lambda m: m.id):
|
||||
if msg.chat_id is None:
|
||||
msg.chat_id = chat.id
|
||||
await self._process_message(client, msg)
|
||||
except Exception:
|
||||
logger.exception("max_history_fetch_failed", chat_id=chat.id)
|
||||
@@ -168,25 +188,21 @@ class MaxListener:
|
||||
logger.exception("max_chat_fetch_failed", chat_id=message.chat_id)
|
||||
return
|
||||
|
||||
is_dm = chat.type == ChatType.DIALOG or getattr(chat, "is_dialog", False)
|
||||
chat_title, _ = build_chat_title(chat, self._settings.ls_topic_prefix)
|
||||
chat_meta = extract_chat_meta(chat, self._settings.ls_topic_prefix)
|
||||
chat_icon_url = await resolve_chat_icon_url(
|
||||
client, chat, self._holder.my_user_id
|
||||
)
|
||||
chat_icon_local_path = await self._download_chat_icon(chat_icon_url)
|
||||
|
||||
sender_name: str | None = None
|
||||
if not is_dm and message.sender:
|
||||
if not chat_meta.is_dm and message.sender:
|
||||
try:
|
||||
user = await client.get_user(message.sender)
|
||||
sender_name = resolve_sender_name(user)
|
||||
except Exception:
|
||||
sender_name = f"User {message.sender}"
|
||||
|
||||
reply_to: int | None = None
|
||||
if message.options and isinstance(message.options, dict):
|
||||
reply_to = message.options.get("replyTo")
|
||||
if reply_to is None and message.prev_message_id:
|
||||
try:
|
||||
reply_to = int(message.prev_message_id)
|
||||
except (TypeError, ValueError):
|
||||
reply_to = None
|
||||
reply_to = get_reply_target(message)
|
||||
|
||||
try:
|
||||
if forwarded_attaches is not None:
|
||||
@@ -213,7 +229,8 @@ class MaxListener:
|
||||
is_forward=forwarded is not None,
|
||||
text_len=len(effective_text),
|
||||
media_count=len(media),
|
||||
is_dm=is_dm,
|
||||
is_dm=chat_meta.is_dm,
|
||||
reply_to=reply_to,
|
||||
)
|
||||
|
||||
incoming = MaxIncomingMessage(
|
||||
@@ -222,8 +239,13 @@ class MaxListener:
|
||||
text=effective_text,
|
||||
sender_id=message.sender,
|
||||
sender_name=sender_name,
|
||||
is_dm=is_dm,
|
||||
chat_title=chat_title,
|
||||
is_dm=chat_meta.is_dm,
|
||||
chat_title=chat_meta.topic_title,
|
||||
chat_name=chat_meta.chat_name,
|
||||
chat_icon_url=chat_icon_url,
|
||||
chat_icon_local_path=chat_icon_local_path,
|
||||
participants_count=chat_meta.participants_count,
|
||||
max_chat_link=chat_meta.link,
|
||||
reply_to_max_message_id=reply_to,
|
||||
media=media,
|
||||
)
|
||||
|
||||
@@ -44,6 +44,43 @@ def resolve_file_name(item: MediaItem | dict) -> str:
|
||||
return _default_file_name(kind)
|
||||
|
||||
|
||||
def max_image_url_candidates(url: str) -> list[str]:
|
||||
normalized = url.strip()
|
||||
if not normalized:
|
||||
return []
|
||||
if normalized.startswith("//"):
|
||||
normalized = f"https:{normalized}"
|
||||
candidates = [normalized]
|
||||
if "i.oneme.ru" in normalized and "size=" not in normalized:
|
||||
sep = "&" if "?" in normalized else "?"
|
||||
for size in (512, 256, 128):
|
||||
sized = f"{normalized}{sep}size={size}"
|
||||
if sized not in candidates:
|
||||
candidates.append(sized)
|
||||
return candidates
|
||||
|
||||
|
||||
async def download_url_to_file(
|
||||
session: aiohttp.ClientSession,
|
||||
url: str,
|
||||
dest: Path,
|
||||
) -> bool:
|
||||
return await _download_url(session, url, dest)
|
||||
|
||||
|
||||
async def download_max_image_url(
|
||||
session: aiohttp.ClientSession,
|
||||
url: str,
|
||||
dest: Path,
|
||||
) -> bool:
|
||||
for candidate in max_image_url_candidates(url):
|
||||
if await _download_url(session, candidate, dest):
|
||||
if dest.stat().st_size > 0:
|
||||
return True
|
||||
dest.unlink(missing_ok=True)
|
||||
return False
|
||||
|
||||
|
||||
async def _download_url(
|
||||
session: aiohttp.ClientSession,
|
||||
url: str,
|
||||
|
||||
@@ -10,6 +10,16 @@ class ChatMapping:
|
||||
is_dm: bool
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class MaxChatMeta:
|
||||
topic_title: str
|
||||
chat_name: str
|
||||
is_dm: bool
|
||||
icon_url: str | None
|
||||
participants_count: int
|
||||
link: str | None
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class MaxIncomingMessage:
|
||||
max_chat_id: int
|
||||
@@ -19,6 +29,11 @@ class MaxIncomingMessage:
|
||||
sender_name: str | None
|
||||
is_dm: bool
|
||||
chat_title: str
|
||||
chat_name: str
|
||||
chat_icon_url: str | None
|
||||
chat_icon_local_path: str | None
|
||||
participants_count: int
|
||||
max_chat_link: str | None
|
||||
reply_to_max_message_id: int | None
|
||||
media: list[dict]
|
||||
|
||||
|
||||
@@ -33,6 +33,11 @@ class Max2TgTask(BaseModel):
|
||||
sender_name: str | None = None
|
||||
is_dm: bool = False
|
||||
chat_title: str
|
||||
chat_name: str = ""
|
||||
chat_icon_url: str | None = None
|
||||
chat_icon_local_path: str | None = None
|
||||
participants_count: int = 0
|
||||
max_chat_link: str | None = None
|
||||
needs_new_topic: bool = False
|
||||
reply_to_tg_message_id: int | None = None
|
||||
media: list[MediaItem] = Field(default_factory=list)
|
||||
|
||||
@@ -50,6 +50,21 @@ class MessageRouter:
|
||||
max_chat_id=message.max_chat_id,
|
||||
thread_id=mapping.tg_thread_id,
|
||||
)
|
||||
else:
|
||||
all_mappings = await self._storage.list_mappings()
|
||||
logger.info(
|
||||
"router_no_mapping",
|
||||
max_chat_id=message.max_chat_id,
|
||||
reason="chat_mapping_not_found_in_storage",
|
||||
will_create_topic=True,
|
||||
total_mappings_in_storage=len(all_mappings),
|
||||
known_max_chat_ids=[m.max_chat_id for m in all_mappings] if all_mappings else [],
|
||||
hint=(
|
||||
"empty_db_after_container_recreate"
|
||||
if not all_mappings
|
||||
else "mapping_missing_for_this_chat_only"
|
||||
),
|
||||
)
|
||||
|
||||
reply_to_tg: int | None = None
|
||||
if message.reply_to_max_message_id is not None:
|
||||
@@ -79,6 +94,11 @@ class MessageRouter:
|
||||
sender_name=message.sender_name,
|
||||
is_dm=message.is_dm,
|
||||
chat_title=message.chat_title,
|
||||
chat_name=message.chat_name,
|
||||
chat_icon_url=message.chat_icon_url,
|
||||
chat_icon_local_path=message.chat_icon_local_path,
|
||||
participants_count=message.participants_count,
|
||||
max_chat_link=message.max_chat_link,
|
||||
needs_new_topic=needs_new_topic,
|
||||
reply_to_tg_message_id=reply_to_tg,
|
||||
media=media,
|
||||
|
||||
@@ -2,6 +2,17 @@ from aiogram.types import Message
|
||||
|
||||
from app.models.domain import TgIncomingMessage
|
||||
|
||||
TG_CAPTION_MAX_LEN = 1024
|
||||
TG_MESSAGE_MAX_LEN = 4096
|
||||
|
||||
|
||||
def split_text_for_media_caption(text: str) -> tuple[str | None, str | None]:
|
||||
if not text:
|
||||
return None, None
|
||||
if len(text) <= TG_CAPTION_MAX_LEN:
|
||||
return text, None
|
||||
return None, text
|
||||
|
||||
|
||||
def format_tg_author(message: Message) -> tuple[str, str | None]:
|
||||
user = message.from_user
|
||||
@@ -12,6 +23,25 @@ def format_tg_author(message: Message) -> tuple[str, str | None]:
|
||||
return name, user.username
|
||||
|
||||
|
||||
def format_topic_pin_text(
|
||||
*,
|
||||
chat_name: str,
|
||||
is_dm: bool,
|
||||
max_chat_id: int,
|
||||
participants_count: int,
|
||||
max_chat_link: str | None,
|
||||
) -> str:
|
||||
chat_type = "PRIVATE" if is_dm else "CHAT"
|
||||
lines = [
|
||||
f"{chat_name} · Тип: {chat_type}",
|
||||
f"id: {max_chat_id}",
|
||||
f"Участников: {participants_count}",
|
||||
]
|
||||
if max_chat_link:
|
||||
lines.extend(["", f"🔗 {max_chat_link}"])
|
||||
return "\n".join(lines)
|
||||
|
||||
|
||||
def format_tg_to_max_text(author_name: str, username: str | None, text: str) -> str:
|
||||
handle = f" (@{username})" if username else ""
|
||||
return f"{author_name}{handle}:\n{text}"
|
||||
|
||||
+222
-32
@@ -1,7 +1,10 @@
|
||||
import asyncio
|
||||
from pathlib import Path
|
||||
from uuid import uuid4
|
||||
|
||||
import aiohttp
|
||||
import structlog
|
||||
from aiogram.exceptions import TelegramBadRequest
|
||||
from aiogram.types import (
|
||||
FSInputFile,
|
||||
InputMediaDocument,
|
||||
@@ -11,7 +14,13 @@ from aiogram.types import (
|
||||
)
|
||||
|
||||
from app.config import Settings
|
||||
from app.media_transfer import cleanup_paths, download_media_item, resolve_file_name, tmp_dir
|
||||
from app.media_transfer import (
|
||||
cleanup_paths,
|
||||
download_max_image_url,
|
||||
download_media_item,
|
||||
resolve_file_name,
|
||||
tmp_dir,
|
||||
)
|
||||
from app.max_layer.client_holder import MaxClientHolder
|
||||
from app.max_layer.formatter import format_max_text
|
||||
from app.models.domain import ChatMapping, MaxIncomingMessage
|
||||
@@ -24,11 +33,19 @@ from app.models.tasks import (
|
||||
from app.queue.protocols import QueuePort
|
||||
from app.storage.protocols import StoragePort
|
||||
from app.telegram_layer.bot_holder import BotHolder
|
||||
from app.telegram_layer.formatter import (
|
||||
format_topic_pin_text,
|
||||
split_text_for_media_caption,
|
||||
)
|
||||
from app.topic_locks import TopicLockRegistry
|
||||
|
||||
logger = structlog.get_logger(__name__)
|
||||
|
||||
|
||||
def _is_thread_not_found(exc: TelegramBadRequest) -> bool:
|
||||
return "message thread not found" in exc.message.lower()
|
||||
|
||||
|
||||
class TelegramWorker:
|
||||
def __init__(
|
||||
self,
|
||||
@@ -73,6 +90,117 @@ class TelegramWorker:
|
||||
)
|
||||
await asyncio.sleep(self._settings.tg_rate_limit_delay_sec)
|
||||
|
||||
async def _create_and_save_topic(
|
||||
self,
|
||||
bot,
|
||||
task: Max2TgTask,
|
||||
*,
|
||||
reason: str,
|
||||
old_thread_id: int | None = None,
|
||||
) -> ChatMapping:
|
||||
topic = await bot.create_forum_topic(
|
||||
chat_id=self._settings.tg_forum_channel_id,
|
||||
name=task.chat_title[:128],
|
||||
)
|
||||
mapping = ChatMapping(
|
||||
max_chat_id=task.max_chat_id,
|
||||
tg_chat_id=self._settings.tg_forum_channel_id,
|
||||
tg_thread_id=topic.message_thread_id,
|
||||
display_name=task.chat_title,
|
||||
is_dm=task.is_dm,
|
||||
)
|
||||
await self._storage.save_mapping(mapping)
|
||||
await self._send_and_pin_topic_info(bot, topic.message_thread_id, task)
|
||||
logger.info(
|
||||
"tg_topic_created",
|
||||
max_chat_id=task.max_chat_id,
|
||||
thread_id=topic.message_thread_id,
|
||||
title=task.chat_title,
|
||||
reason=reason,
|
||||
old_thread_id=old_thread_id,
|
||||
)
|
||||
return mapping
|
||||
|
||||
async def _send_and_pin_topic_info(
|
||||
self,
|
||||
bot,
|
||||
thread_id: int,
|
||||
task: Max2TgTask,
|
||||
) -> None:
|
||||
chat_name = task.chat_name or task.chat_title
|
||||
caption = format_topic_pin_text(
|
||||
chat_name=chat_name,
|
||||
is_dm=task.is_dm,
|
||||
max_chat_id=task.max_chat_id,
|
||||
participants_count=task.participants_count,
|
||||
max_chat_link=task.max_chat_link,
|
||||
)
|
||||
kwargs: dict = {
|
||||
"chat_id": self._settings.tg_forum_channel_id,
|
||||
"message_thread_id": thread_id,
|
||||
}
|
||||
|
||||
downloaded: Path | None = None
|
||||
icon_path: Path | None = None
|
||||
try:
|
||||
icon_path = await self._resolve_topic_icon_path(task)
|
||||
if icon_path is not None:
|
||||
sent = await bot.send_photo(
|
||||
photo=FSInputFile(icon_path),
|
||||
caption=caption,
|
||||
**kwargs,
|
||||
)
|
||||
if (
|
||||
not task.chat_icon_local_path
|
||||
or icon_path.resolve() != Path(task.chat_icon_local_path).resolve()
|
||||
):
|
||||
downloaded = icon_path
|
||||
else:
|
||||
sent = await bot.send_message(text=caption, **kwargs)
|
||||
|
||||
await bot.pin_chat_message(
|
||||
chat_id=self._settings.tg_forum_channel_id,
|
||||
message_id=sent.message_id,
|
||||
disable_notification=True,
|
||||
)
|
||||
logger.info(
|
||||
"tg_topic_info_pinned",
|
||||
max_chat_id=task.max_chat_id,
|
||||
thread_id=thread_id,
|
||||
message_id=sent.message_id,
|
||||
has_icon=icon_path is not None,
|
||||
)
|
||||
except Exception:
|
||||
logger.exception(
|
||||
"tg_topic_info_pin_failed",
|
||||
max_chat_id=task.max_chat_id,
|
||||
thread_id=thread_id,
|
||||
)
|
||||
finally:
|
||||
if downloaded is not None:
|
||||
cleanup_paths([downloaded])
|
||||
|
||||
async def _resolve_topic_icon_path(self, task: Max2TgTask) -> Path | None:
|
||||
if task.chat_icon_local_path:
|
||||
local = Path(task.chat_icon_local_path)
|
||||
if local.is_file() and local.stat().st_size > 0:
|
||||
return local
|
||||
|
||||
if not task.chat_icon_url:
|
||||
return None
|
||||
|
||||
dest = self._tmp_dir / f"topic_icon_{task.max_chat_id}_{uuid4().hex}.jpg"
|
||||
async with aiohttp.ClientSession() as session:
|
||||
if await download_max_image_url(session, task.chat_icon_url, dest):
|
||||
return dest
|
||||
logger.warning(
|
||||
"tg_topic_icon_download_failed",
|
||||
max_chat_id=task.max_chat_id,
|
||||
url=task.chat_icon_url,
|
||||
)
|
||||
dest.unlink(missing_ok=True)
|
||||
return None
|
||||
|
||||
async def _send_to_tg(self, task: Max2TgTask) -> None:
|
||||
bot = await self._bot_holder.wait_bot()
|
||||
logger.debug(
|
||||
@@ -86,23 +214,25 @@ class TelegramWorker:
|
||||
async with self._topic_locks.lock(task.max_chat_id):
|
||||
mapping = await self._storage.get_mapping_by_max_chat(task.max_chat_id)
|
||||
if mapping is None:
|
||||
topic = await bot.create_forum_topic(
|
||||
chat_id=self._settings.tg_forum_channel_id,
|
||||
name=task.chat_title[:128],
|
||||
)
|
||||
mapping = ChatMapping(
|
||||
max_chat_id=task.max_chat_id,
|
||||
tg_chat_id=self._settings.tg_forum_channel_id,
|
||||
tg_thread_id=topic.message_thread_id,
|
||||
display_name=task.chat_title,
|
||||
is_dm=task.is_dm,
|
||||
)
|
||||
await self._storage.save_mapping(mapping)
|
||||
logger.info(
|
||||
"tg_topic_created",
|
||||
"tg_topic_creating",
|
||||
max_chat_id=task.max_chat_id,
|
||||
thread_id=topic.message_thread_id,
|
||||
title=task.chat_title,
|
||||
reason="no_mapping_in_storage_at_send_time",
|
||||
needs_new_topic_from_router=task.needs_new_topic,
|
||||
sqlite_path=self._settings.sqlite_path,
|
||||
hint="mapping_missing_in_sqlite_will_call_create_forum_topic",
|
||||
)
|
||||
mapping = await self._create_and_save_topic(
|
||||
bot, task, reason="no_mapping_in_storage"
|
||||
)
|
||||
elif task.needs_new_topic:
|
||||
logger.info(
|
||||
"tg_topic_reused",
|
||||
max_chat_id=task.max_chat_id,
|
||||
thread_id=mapping.tg_thread_id,
|
||||
reason="mapping_found_despite_needs_new_topic_flag",
|
||||
hint="mapping_appeared_between_router_enqueue_and_worker_send",
|
||||
)
|
||||
else:
|
||||
logger.debug(
|
||||
@@ -112,7 +242,6 @@ class TelegramWorker:
|
||||
)
|
||||
|
||||
assert mapping is not None
|
||||
thread_id = mapping.tg_thread_id
|
||||
|
||||
formatted = format_max_text(
|
||||
MaxIncomingMessage(
|
||||
@@ -123,13 +252,41 @@ class TelegramWorker:
|
||||
sender_name=task.sender_name,
|
||||
is_dm=task.is_dm,
|
||||
chat_title=task.chat_title,
|
||||
chat_name=task.chat_name,
|
||||
chat_icon_url=task.chat_icon_url,
|
||||
chat_icon_local_path=task.chat_icon_local_path,
|
||||
participants_count=task.participants_count,
|
||||
max_chat_link=task.max_chat_link,
|
||||
reply_to_max_message_id=None,
|
||||
media=[],
|
||||
)
|
||||
)
|
||||
|
||||
reply_to = task.reply_to_tg_message_id
|
||||
try:
|
||||
sent = await self._dispatch_content(
|
||||
bot, thread_id, formatted, task.media, task.reply_to_tg_message_id
|
||||
bot, mapping.tg_thread_id, formatted, task.media, reply_to
|
||||
)
|
||||
except TelegramBadRequest as exc:
|
||||
if not _is_thread_not_found(exc):
|
||||
raise
|
||||
old_thread_id = mapping.tg_thread_id
|
||||
logger.info(
|
||||
"tg_topic_stale",
|
||||
max_chat_id=task.max_chat_id,
|
||||
old_thread_id=old_thread_id,
|
||||
reason="telegram_message_thread_not_found",
|
||||
action="recreate_topic_and_update_binding",
|
||||
)
|
||||
async with self._topic_locks.lock(task.max_chat_id):
|
||||
mapping = await self._create_and_save_topic(
|
||||
bot,
|
||||
task,
|
||||
reason="telegram_thread_recreated",
|
||||
old_thread_id=old_thread_id,
|
||||
)
|
||||
sent = await self._dispatch_content(
|
||||
bot, mapping.tg_thread_id, formatted, task.media, None
|
||||
)
|
||||
if sent is None:
|
||||
raise RuntimeError("Telegram API returned no message")
|
||||
@@ -145,7 +302,27 @@ class TelegramWorker:
|
||||
"tg_message_sent",
|
||||
max_message_id=task.max_message_id,
|
||||
tg_message_id=sent.message_id,
|
||||
thread_id=mapping.tg_thread_id,
|
||||
)
|
||||
|
||||
async def _send_overflow_text(
|
||||
self,
|
||||
bot,
|
||||
*,
|
||||
thread_id: int,
|
||||
text: str,
|
||||
reply_to_message_id: int,
|
||||
) -> None:
|
||||
await bot.send_message(
|
||||
chat_id=self._settings.tg_forum_channel_id,
|
||||
message_thread_id=thread_id,
|
||||
text=text,
|
||||
reply_to_message_id=reply_to_message_id,
|
||||
)
|
||||
logger.info(
|
||||
"tg_caption_overflow_sent",
|
||||
thread_id=thread_id,
|
||||
text_len=len(text),
|
||||
)
|
||||
|
||||
async def _dispatch_content(
|
||||
@@ -157,15 +334,18 @@ class TelegramWorker:
|
||||
reply_to: int | None,
|
||||
):
|
||||
chat_id = self._settings.tg_forum_channel_id
|
||||
kwargs = {
|
||||
kwargs: dict = {
|
||||
"chat_id": chat_id,
|
||||
"message_thread_id": thread_id,
|
||||
"reply_to_message_id": reply_to,
|
||||
}
|
||||
if reply_to is not None:
|
||||
kwargs["reply_to_message_id"] = reply_to
|
||||
|
||||
if not media:
|
||||
return await bot.send_message(text=text or " ", **kwargs)
|
||||
|
||||
caption, overflow = split_text_for_media_caption(text)
|
||||
|
||||
downloaded: list = []
|
||||
try:
|
||||
max_client = self._max_holder.client
|
||||
@@ -199,33 +379,43 @@ class TelegramWorker:
|
||||
if len(uploads) == 1:
|
||||
item, upload = uploads[0]
|
||||
if item.kind == "photo":
|
||||
return await bot.send_photo(
|
||||
photo=upload, caption=text or None, **kwargs
|
||||
sent = await bot.send_photo(
|
||||
photo=upload, caption=caption, **kwargs
|
||||
)
|
||||
if item.kind == "video":
|
||||
return await bot.send_video(
|
||||
video=upload, caption=text or None, **kwargs
|
||||
elif item.kind == "video":
|
||||
sent = await bot.send_video(
|
||||
video=upload, caption=caption, **kwargs
|
||||
)
|
||||
return await bot.send_document(
|
||||
document=upload, caption=text or None, **kwargs
|
||||
else:
|
||||
sent = await bot.send_document(
|
||||
document=upload, caption=caption, **kwargs
|
||||
)
|
||||
|
||||
else:
|
||||
group = []
|
||||
for idx, (item, upload) in enumerate(uploads):
|
||||
caption = text if idx == 0 else None
|
||||
item_caption = caption if idx == 0 else None
|
||||
if item.kind == "photo":
|
||||
group.append(InputMediaPhoto(media=upload, caption=caption))
|
||||
group.append(InputMediaPhoto(media=upload, caption=item_caption))
|
||||
elif item.kind == "video":
|
||||
group.append(InputMediaVideo(media=upload, caption=caption))
|
||||
group.append(InputMediaVideo(media=upload, caption=item_caption))
|
||||
else:
|
||||
group.append(InputMediaDocument(media=upload, caption=caption))
|
||||
group.append(InputMediaDocument(media=upload, caption=item_caption))
|
||||
|
||||
messages = await bot.send_media_group(
|
||||
chat_id=chat_id,
|
||||
message_thread_id=thread_id,
|
||||
media=group,
|
||||
)
|
||||
return messages[0]
|
||||
sent = messages[0]
|
||||
|
||||
if overflow:
|
||||
await self._send_overflow_text(
|
||||
bot,
|
||||
thread_id=thread_id,
|
||||
text=overflow,
|
||||
reply_to_message_id=sent.message_id,
|
||||
)
|
||||
return sent
|
||||
finally:
|
||||
cleanup_paths(downloaded)
|
||||
|
||||
|
||||
+1
-1
@@ -19,7 +19,7 @@ services:
|
||||
- TG_BOT_TOKEN=${TG_BOT_TOKEN}
|
||||
- TG_FORUM_CHANNEL_ID=${TG_FORUM_CHANNEL_ID}
|
||||
- FALLBACK_USER_ID=${FALLBACK_USER_ID}
|
||||
- DATABASE_URL=sqlite+aiosqlite:///app/data/bridge.db
|
||||
- DATABASE_URL=sqlite+aiosqlite:///data/bridge.db
|
||||
- REDIS_URL=redis://redis:6379/0
|
||||
- TG_RATE_LIMIT_DELAY_SEC=3.5
|
||||
- MAX_RATE_LIMIT_DELAY_SEC=1.0
|
||||
|
||||
+2
-2
@@ -122,7 +122,7 @@ services:
|
||||
- TG_BOT_TOKEN=${TG_BOT_TOKEN}
|
||||
- TG_FORUM_CHANNEL_ID=${TG_FORUM_CHANNEL_ID}
|
||||
- FALLBACK_USER_ID=${FALLBACK_USER_ID}
|
||||
- DATABASE_URL=sqlite+aiosqlite:///app/data/bridge.db
|
||||
- DATABASE_URL=sqlite+aiosqlite:///data/bridge.db
|
||||
- REDIS_URL=redis://redis:6379/0
|
||||
- TG_RATE_LIMIT_DELAY_SEC=3.5
|
||||
- MAX_RATE_LIMIT_DELAY_SEC=1.0
|
||||
@@ -211,7 +211,7 @@ volumes:
|
||||
| `TG_BOT_TOKEN` | Токен Telegram-бота от @BotFather | `123456:ABC-DEF1234...` |
|
||||
| `TG_FORUM_CHANNEL_ID` | ID Форум-канала (строго с `-100`) | `-1001234567890` |
|
||||
| `FALLBACK_USER_ID` | Telegram ID администратора | `987654321` |
|
||||
| `DATABASE_URL` | Строка подключения к SQLite | `sqlite+aiosqlite:///app/data/bridge.db` |
|
||||
| `DATABASE_URL` | Строка подключения к SQLite (относительно `WORKDIR=/app`) | `sqlite+aiosqlite:///data/bridge.db` |
|
||||
| `REDIS_URL` | Строка подключения к Redis | `redis://redis:6379/0` |
|
||||
| `TG_RATE_LIMIT_DELAY_SEC` | Задержка между отправками в TG | `3.5` |
|
||||
| `MAX_RATE_LIMIT_DELAY_SEC` | Задержка между отправками в MAX | `1.0` |
|
||||
|
||||
Reference in New Issue
Block a user