From aecfe4dedfed5561a0b20751d319e6d4672d4b80 Mon Sep 17 00:00:00 2001 From: kislovdm Date: Fri, 12 Jun 2026 01:10:44 +0300 Subject: [PATCH] +++ --- .env.example | 8 +-- app/config.py | 11 ++-- app/main.py | 26 +++++++++ app/max_layer/formatter.py | 15 ++++++ app/max_layer/listener.py | 11 ++-- app/router/router.py | 15 ++++++ app/telegram_layer/worker.py | 102 +++++++++++++++++++++++++++-------- docker-compose.yml | 2 +- tech-specs.md | 4 +- 9 files changed, 151 insertions(+), 43 deletions(-) diff --git a/.env.example b/.env.example index 2534065..0d76a5f 100644 --- a/.env.example +++ b/.env.example @@ -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 diff --git a/app/config.py b/app/config.py index 755d2c0..66f9f73 100644 --- a/app/config.py +++ b/app/config.py @@ -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": diff --git a/app/main.py b/app/main.py index 150fb6d..6a2de3e 100644 --- a/app/main.py +++ b/app/main.py @@ -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() diff --git a/app/max_layer/formatter.py b/app/max_layer/formatter.py index 7e76ec2..bdf57a9 100644 --- a/app/max_layer/formatter.py +++ b/app/max_layer/formatter.py @@ -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: diff --git a/app/max_layer/listener.py b/app/max_layer/listener.py index 143eef8..c2969d1 100644 --- a/app/max_layer/listener.py +++ b/app/max_layer/listener.py @@ -12,6 +12,7 @@ from app.max_layer.formatter import ( build_chat_title, extract_forwarded_content, format_forwarded_text, + get_reply_target, resolve_media, resolve_raw_attaches, resolve_sender_name, @@ -179,14 +180,7 @@ class MaxListener: 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: @@ -214,6 +208,7 @@ class MaxListener: text_len=len(effective_text), media_count=len(media), is_dm=is_dm, + reply_to=reply_to, ) incoming = MaxIncomingMessage( diff --git a/app/router/router.py b/app/router/router.py index 88379d4..fa0740e 100644 --- a/app/router/router.py +++ b/app/router/router.py @@ -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: diff --git a/app/telegram_layer/worker.py b/app/telegram_layer/worker.py index 6a88756..4df85e2 100644 --- a/app/telegram_layer/worker.py +++ b/app/telegram_layer/worker.py @@ -2,6 +2,7 @@ import asyncio import aiohttp import structlog +from aiogram.exceptions import TelegramBadRequest from aiogram.types import ( FSInputFile, InputMediaDocument, @@ -29,6 +30,10 @@ 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 +78,36 @@ 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) + 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_to_tg(self, task: Max2TgTask) -> None: bot = await self._bot_holder.wait_bot() logger.debug( @@ -86,23 +121,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 +149,6 @@ class TelegramWorker: ) assert mapping is not None - thread_id = mapping.tg_thread_id formatted = format_max_text( MaxIncomingMessage( @@ -128,9 +164,32 @@ class TelegramWorker: ) ) - sent = await self._dispatch_content( - bot, thread_id, formatted, task.media, task.reply_to_tg_message_id - ) + reply_to = task.reply_to_tg_message_id + try: + sent = await self._dispatch_content( + 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 +204,7 @@ class TelegramWorker: "tg_message_sent", max_message_id=task.max_message_id, tg_message_id=sent.message_id, - thread_id=thread_id, + thread_id=mapping.tg_thread_id, ) async def _dispatch_content( @@ -157,11 +216,12 @@ 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) diff --git a/docker-compose.yml b/docker-compose.yml index c1382f5..6ff9056 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -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 diff --git a/tech-specs.md b/tech-specs.md index ea921e8..e960a01 100644 --- a/tech-specs.md +++ b/tech-specs.md @@ -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` |