ИСправление ошибки с двойным поллингом
This commit is contained in:
@@ -93,7 +93,9 @@ class TelegramToMaxBridge:
|
|||||||
if self._health:
|
if self._health:
|
||||||
self._health.mark_telegram_error()
|
self._health.mark_telegram_error()
|
||||||
logger.exception("Telegram polling loop error")
|
logger.exception("Telegram polling loop error")
|
||||||
await asyncio.sleep(2)
|
# 409 Conflict: где-то еще идет getUpdates (другой инстанс или webhook/второй poller).
|
||||||
|
# Делаем backoff, чтобы не долбить API.
|
||||||
|
await asyncio.sleep(10)
|
||||||
|
|
||||||
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
|
||||||
|
|||||||
+21
-18
@@ -108,7 +108,19 @@ class TelegramClient:
|
|||||||
result = data.get("result", [])
|
result = data.get("result", [])
|
||||||
if not isinstance(result, list):
|
if not isinstance(result, list):
|
||||||
return []
|
return []
|
||||||
return [u for u in result if isinstance(u, dict)]
|
updates = [u for u in result if isinstance(u, dict)]
|
||||||
|
# Важно: не делаем getUpdates нигде больше (иначе 409 Conflict).
|
||||||
|
# Наполняем кэш чатов только из этого потока.
|
||||||
|
for upd in updates:
|
||||||
|
for container in ("message", "edited_message", "channel_post", "edited_channel_post"):
|
||||||
|
msg = upd.get(container)
|
||||||
|
if not isinstance(msg, dict):
|
||||||
|
continue
|
||||||
|
chat = msg.get("chat")
|
||||||
|
if not isinstance(chat, dict):
|
||||||
|
continue
|
||||||
|
self._cache_chat(chat)
|
||||||
|
return updates
|
||||||
|
|
||||||
async def get_file_url(self, file_id: str) -> str:
|
async def get_file_url(self, file_id: str) -> str:
|
||||||
data = await self._request("getFile", {"file_id": file_id})
|
data = await self._request("getFile", {"file_id": file_id})
|
||||||
@@ -131,6 +143,12 @@ class TelegramClient:
|
|||||||
},
|
},
|
||||||
)
|
)
|
||||||
|
|
||||||
|
def _cache_chat(self, chat: dict[str, Any]) -> None:
|
||||||
|
title_value = self._extract_chat_title(chat)
|
||||||
|
chat_id = chat.get("id")
|
||||||
|
if title_value and chat_id is not None:
|
||||||
|
self._chat_title_to_id[self._normalize_title(title_value)] = str(chat_id)
|
||||||
|
|
||||||
async def _find_chat_id_by_title(self, chat_title: str) -> str | None:
|
async def _find_chat_id_by_title(self, chat_title: str) -> str | None:
|
||||||
normalized = self._normalize_title(chat_title)
|
normalized = self._normalize_title(chat_title)
|
||||||
if not normalized:
|
if not normalized:
|
||||||
@@ -139,23 +157,8 @@ class TelegramClient:
|
|||||||
cached = self._chat_title_to_id.get(normalized)
|
cached = self._chat_title_to_id.get(normalized)
|
||||||
if cached:
|
if cached:
|
||||||
return cached
|
return cached
|
||||||
|
# Не дергаем getUpdates здесь — это вызовет конфликт с polling циклом.
|
||||||
response = await self._request("getUpdates", {"timeout": 0, "limit": 100})
|
return None
|
||||||
for update in response.get("result", []):
|
|
||||||
for container in ("message", "edited_message", "channel_post", "edited_channel_post"):
|
|
||||||
message = update.get(container)
|
|
||||||
if not isinstance(message, dict):
|
|
||||||
continue
|
|
||||||
chat = message.get("chat")
|
|
||||||
if not isinstance(chat, dict):
|
|
||||||
continue
|
|
||||||
|
|
||||||
title_value = self._extract_chat_title(chat)
|
|
||||||
chat_id = chat.get("id")
|
|
||||||
if title_value and chat_id is not None:
|
|
||||||
self._chat_title_to_id[self._normalize_title(title_value)] = str(chat_id)
|
|
||||||
|
|
||||||
return self._chat_title_to_id.get(normalized)
|
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def _extract_chat_title(chat: dict[str, Any]) -> str:
|
def _extract_chat_title(chat: dict[str, Any]) -> str:
|
||||||
|
|||||||
Reference in New Issue
Block a user