From b0f4f1b202bdd16b9a76f65f29632259773082a7 Mon Sep 17 00:00:00 2001 From: kislovdm Date: Fri, 10 Apr 2026 10:51:23 +0300 Subject: [PATCH] =?UTF-8?q?=D0=B8=D1=81=D0=BF=D1=80=D0=B0=D0=B2=D0=BB?= =?UTF-8?q?=D1=8F=D0=B5=D0=BC=20=D0=BF=D1=80=D0=BE=D0=B1=D0=BB=D0=B5=D0=BC?= =?UTF-8?q?=D1=83=20=D1=81=20=D0=B1=D0=BB=D0=BE=D0=BA=D0=B8=D1=80=D0=BE?= =?UTF-8?q?=D0=B2=D0=BA=D0=BE=D0=B9=20=D1=82=D0=B3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/main.py | 12 +++++++++-- src/reverse_bridge.py | 50 +++++++++++++++++++++++++++---------------- 2 files changed, 41 insertions(+), 21 deletions(-) diff --git a/src/main.py b/src/main.py index 4483405..4520071 100644 --- a/src/main.py +++ b/src/main.py @@ -51,13 +51,21 @@ def main() -> None: reverse_bridge = TelegramToMaxBridge(max_client=max_client, telegram=telegram_client, storage=storage, health=health) logger = logging.getLogger("max2telegram") + reverse_bridge_task: asyncio.Task[None] | None = None + max_probe_task: asyncio.Task[None] | None = None @max_client.on_start async def on_start() -> None: + nonlocal reverse_bridge_task, max_probe_task logger.info("Max client started as %s", max_client.me.id) health.mark_max_ok() - asyncio.create_task(reverse_bridge.start()) - asyncio.create_task(_max_probe_loop(max_client=max_client, health=health)) + if reverse_bridge_task is None or reverse_bridge_task.done(): + reverse_bridge_task = asyncio.create_task(reverse_bridge.start(), name="reverse-bridge") + if max_probe_task is None or max_probe_task.done(): + max_probe_task = asyncio.create_task( + _max_probe_loop(max_client=max_client, health=health), + name="max-probe-loop", + ) @max_client.on_message() async def on_message(message: Message) -> None: diff --git a/src/reverse_bridge.py b/src/reverse_bridge.py index da51b21..598996d 100644 --- a/src/reverse_bridge.py +++ b/src/reverse_bridge.py @@ -74,29 +74,41 @@ class TelegramToMaxBridge: self._media_groups: dict[tuple[str, str], _MediaGroupBuffer] = {} self._media_group_grace_sec = 1.2 + self._start_lock = asyncio.Lock() + self._is_running = False async def start(self) -> None: - me = await self._telegram.get_me() - self._bot_id = str(me.get("id") or "") - if not self._bot_id: - raise RuntimeError("Cannot resolve Telegram bot id (getMe)") + async with self._start_lock: + if self._is_running: + logger.warning("Telegram->Max bridge start skipped: poller is already running") + return + self._is_running = True - self._refresh_max_chat_cache() - logger.info("Telegram->Max bridge started (bot_id=%s)", self._bot_id) + try: + me = await self._telegram.get_me() + self._bot_id = str(me.get("id") or "") + if not self._bot_id: + raise RuntimeError("Cannot resolve Telegram bot id (getMe)") - while True: - try: - updates = await self._telegram.get_updates(offset=self._offset, timeout=25, limit=100) - if self._health: - self._health.mark_telegram_ok() - await self._handle_updates(updates) - except Exception: - if self._health: - self._health.mark_telegram_error() - logger.exception("Telegram polling loop error") - # 409 Conflict: где-то еще идет getUpdates (другой инстанс или webhook/второй poller). - # Делаем backoff, чтобы не долбить API. - await asyncio.sleep(10) + self._refresh_max_chat_cache() + logger.info("Telegram->Max bridge started (bot_id=%s)", self._bot_id) + + while True: + try: + updates = await self._telegram.get_updates(offset=self._offset, timeout=25, limit=100) + if self._health: + self._health.mark_telegram_ok() + await self._handle_updates(updates) + except Exception: + if self._health: + self._health.mark_telegram_error() + logger.exception("Telegram polling loop error") + # 409 Conflict: где-то еще идет getUpdates (другой инстанс или webhook/второй poller). + # Делаем backoff, чтобы не долбить API. + await asyncio.sleep(10) + finally: + async with self._start_lock: + self._is_running = False async def _handle_updates(self, updates: list[dict[str, Any]]) -> None: max_update_id = None