исправляем проблему с блокировкой тг
This commit is contained in:
+10
-2
@@ -51,13 +51,21 @@ def main() -> None:
|
|||||||
|
|
||||||
reverse_bridge = TelegramToMaxBridge(max_client=max_client, telegram=telegram_client, storage=storage, health=health)
|
reverse_bridge = TelegramToMaxBridge(max_client=max_client, telegram=telegram_client, storage=storage, health=health)
|
||||||
logger = logging.getLogger("max2telegram")
|
logger = logging.getLogger("max2telegram")
|
||||||
|
reverse_bridge_task: asyncio.Task[None] | None = None
|
||||||
|
max_probe_task: asyncio.Task[None] | None = None
|
||||||
|
|
||||||
@max_client.on_start
|
@max_client.on_start
|
||||||
async def on_start() -> None:
|
async def on_start() -> None:
|
||||||
|
nonlocal reverse_bridge_task, max_probe_task
|
||||||
logger.info("Max client started as %s", max_client.me.id)
|
logger.info("Max client started as %s", max_client.me.id)
|
||||||
health.mark_max_ok()
|
health.mark_max_ok()
|
||||||
asyncio.create_task(reverse_bridge.start())
|
if reverse_bridge_task is None or reverse_bridge_task.done():
|
||||||
asyncio.create_task(_max_probe_loop(max_client=max_client, health=health))
|
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()
|
@max_client.on_message()
|
||||||
async def on_message(message: Message) -> None:
|
async def on_message(message: Message) -> None:
|
||||||
|
|||||||
+31
-19
@@ -74,29 +74,41 @@ class TelegramToMaxBridge:
|
|||||||
|
|
||||||
self._media_groups: dict[tuple[str, str], _MediaGroupBuffer] = {}
|
self._media_groups: dict[tuple[str, str], _MediaGroupBuffer] = {}
|
||||||
self._media_group_grace_sec = 1.2
|
self._media_group_grace_sec = 1.2
|
||||||
|
self._start_lock = asyncio.Lock()
|
||||||
|
self._is_running = False
|
||||||
|
|
||||||
async def start(self) -> None:
|
async def start(self) -> None:
|
||||||
me = await self._telegram.get_me()
|
async with self._start_lock:
|
||||||
self._bot_id = str(me.get("id") or "")
|
if self._is_running:
|
||||||
if not self._bot_id:
|
logger.warning("Telegram->Max bridge start skipped: poller is already running")
|
||||||
raise RuntimeError("Cannot resolve Telegram bot id (getMe)")
|
return
|
||||||
|
self._is_running = True
|
||||||
|
|
||||||
self._refresh_max_chat_cache()
|
try:
|
||||||
logger.info("Telegram->Max bridge started (bot_id=%s)", self._bot_id)
|
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:
|
self._refresh_max_chat_cache()
|
||||||
try:
|
logger.info("Telegram->Max bridge started (bot_id=%s)", self._bot_id)
|
||||||
updates = await self._telegram.get_updates(offset=self._offset, timeout=25, limit=100)
|
|
||||||
if self._health:
|
while True:
|
||||||
self._health.mark_telegram_ok()
|
try:
|
||||||
await self._handle_updates(updates)
|
updates = await self._telegram.get_updates(offset=self._offset, timeout=25, limit=100)
|
||||||
except Exception:
|
if self._health:
|
||||||
if self._health:
|
self._health.mark_telegram_ok()
|
||||||
self._health.mark_telegram_error()
|
await self._handle_updates(updates)
|
||||||
logger.exception("Telegram polling loop error")
|
except Exception:
|
||||||
# 409 Conflict: где-то еще идет getUpdates (другой инстанс или webhook/второй poller).
|
if self._health:
|
||||||
# Делаем backoff, чтобы не долбить API.
|
self._health.mark_telegram_error()
|
||||||
await asyncio.sleep(10)
|
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:
|
async def _handle_updates(self, updates: list[dict[str, Any]]) -> None:
|
||||||
max_update_id = None
|
max_update_id = None
|
||||||
|
|||||||
Reference in New Issue
Block a user