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