gendesign/tradein-mvp/backend/tests/services/tgbot/test_client.py
bot-backend 1fa65eba6b
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 10s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m19s
fix(tg): связь с Telegram не встаёт колом, ответ оператора не теряется
Замер прода за сутки 12.09.2026: 576 строк `network error` в логе `tradein-tgbot`
и 7 полных исчерпаний бюджета ретраев, после которых падала итерация poll loop.
Три причины, все подтверждены на коде и в рантайме.

## Ответ оператора мог пропасть навсегда

`process_update` заканчивался безусловным `finally: save_offset(update_id)`.
Замысел верный — «ядовитый» апдейт не должен блокировать поток, — но он не
отличал неисправимый апдейт от транзиентного сетевого отказа. Оператор отвечает
клиенту в топике, `copy_message` падает по сети, `TelegramNetworkError` улетает
в общий `except Exception`, offset сдвигается. Telegram этот апдейт больше не
отдаст, `record_message` не выполнился, оператор уверен, что ответил. Следа нет
нигде, кроме строчки в логе.

Теперь `process_update` возвращает `bool`. На `TelegramNetworkError` делается
`rollback()`, offset НЕ сохраняется, возвращается `False`, и `run_poll_loop`
прерывает разбор пачки — offset у Telegram единая «высшая отметка», подтверждение
любого следующего апдейта неявно подтвердило бы и этот. Остаток пачки Telegram
отдаст заново.

Переигрывания ограничены сверху `_MAX_NETWORK_REPLAYS = 3`: без потолка «вечно
недоставляемый» апдейт заклинил бы очередь навсегда, а это хуже потери одного
сообщения. На потолке offset всё-таки двигается, но с `logger.error` и с
`chat_id`/`message_id`, по которым человек найдёт ответ в топике и перешлёт
руками. Текст переписки в лог по-прежнему не идёт.

Дубли: `TelegramNetworkError` означает исчерпанный бюджет ретраев, при этом
запрос мог дойти до Telegram, а ответ потеряться. Переигрывание тогда доставит
сообщение второй раз. Это осознанный at-least-once компромисс — дубль видят и
клиент, и оператор, а тихая потеря не видна никому. Полная идемпотентность по
паре (update_id, target_chat_id) потребовала бы новой персистентной таблицы ради
редкого случая; вместо неё число дублей жёстко ограничено сверху.

Ветка `except TelegramApiError` с разбором `error_code == 403` («бот заблокирован»)
не тронута — там повтор действительно ничего не изменит.

## Таймаут задавался скаляром, поэтому connect ждал сорок секунд

`httpx.AsyncClient(timeout=effective_timeout)` разворачивается в
connect=read=write=pool. Для `getUpdates` бюджет ответа 40 секунд (30 держит
Telegram плюс запас), и те же 40 секунд уходили на установку соединения — при
живом connect в 0.036 секунды. Худший цикл: четыре попытки по 40 секунд плюс
backoff, около трёх минут, в течение которых бот не видит ответов оператора.
В логе это ровно те разрывы: 06:40:10, 06:42:22, 06:43:35.

Теперь `httpx.Timeout(connect=5, read=<бюджет вызывающего>, write=10, pool=5)`,
значения в именованных константах. Запас `+10s` у `get_updates` относится к read,
докстринг поправлен.

## Клиент создавался заново на каждую попытку

`httpx.AsyncClient` стоял ВНУТРИ цикла ретраев — keep-alive не было вовсе: полный
TCP+TLS-хендшейк на каждый запрос и на каждый повтор, и заново кидался кубик
«встанет ли коннект». Для long-polling это была основная статья сетевых отказов.
Плюс три HTTP-ручки создавали `TelegramClient` на каждый входящий запрос.

Теперь один ленивый переиспользуемый `AsyncClient` на экземпляр, с `aclose()` и
`async with`. Общий клиент приложения живёт в новом `app/services/tgbot/shared.py`,
создаётся и закрывается в lifespan; воркер бота держит свой на время поллинга.
`keepalive_expiry` задан явно: дефолт httpx — 5 секунд, и с ним пул не давал бы
ничего там, где нужнее всего. Poll loop переиспользует соединение и так, а вот
веб-поддержка шлёт раз в минуты и за 5 секунд теряла бы его каждый раз. Плата за
длинный keep-alive — шанс взять из пула закрытое той стороной соединение; httpx
отдаёт это как `RemoteProtocolError`, который ретраится с #3457.

## Уведомления оператору шли с воркерным бюджетом внутри poll loop

Обе отправки в топик («бот заблокирован», «веб-чат не поддерживает медиа») звались
без своего бюджета, то есть с дефолтом в 5 ретраев и backoff до 30 секунд. Одна
такая отправка стопорила весь цикл на минуты, а её отказ решал судьбу апдейта.
Вынесены в `_notify_topic` с узким бюджетом и собственным `except`: провал
вторичного действия больше не отменяет основную ветку.

## Тесты

`tests/services/tgbot/test_shared.py` — новый, на жизненный цикл общего клиента.
В `test_bridge.py` — сетевой отказ оставляет offset нетронутым и апдейт
переигрывается, потолок разблокирует поток, отказ уведомления не отменяет основную
ветку, прежнее поведение на 403 не изменилось. В `test_client.py` — раздельные
таймауты доезжают до httpx per-request, два вызова используют один `AsyncClient`,
`aclose()` его закрывает.

Прогон по затронутым файлам: 127 passed. Ruff check и format чистые.

Прокси намеренно не добавлялся: замер был на восьми запросах, это не статистика,
и решение инфраструктурное. Если обрывы останутся — мерить сотней попыток отдельно.
2026-09-12 10:13:44 +03:00

411 lines
18 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Unit tests for `app.services.tgbot.client.TelegramClient` retry/backoff logic.
NEVER calls real Telegram API — httpx.MockTransport only (consistent с
tests/services/test_dadata.py). `asyncio.sleep` is patched to a no-op so retry
tests run instantly regardless of configured backoff/retry_after durations.
"""
from __future__ import annotations
import os
from typing import Any
from unittest import mock
import httpx
import pytest
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
from app.services.tgbot.client import (
_CONNECT_TIMEOUT_S,
_POOL_TIMEOUT_S,
_WRITE_TIMEOUT_S,
TelegramApiError,
TelegramClient,
TelegramNetworkError,
)
_REAL_ASYNC_CLIENT = httpx.AsyncClient
def _install_transport(handler) -> list[httpx.AsyncClient]:
"""Подменяет транспорт. Возвращает список СОЗДАННЫХ AsyncClient — по нему
видно, переиспользуется ли один клиент или он плодится на каждый запрос."""
transport = httpx.MockTransport(handler)
created: list[httpx.AsyncClient] = []
def factory(*_: object, **__: object) -> httpx.AsyncClient:
client = _REAL_ASYNC_CLIENT(transport=transport)
created.append(client)
return client
mock.patch("app.services.tgbot.client.httpx.AsyncClient", factory).start()
return created
@pytest.fixture(autouse=True)
def _stop_patches_and_noop_sleep():
sleep_patcher = mock.patch("app.services.tgbot.client.asyncio.sleep", return_value=None)
sleep_patcher.start()
yield
mock.patch.stopall()
async def test_get_updates_happy_path_returns_list() -> None:
def handler(request: httpx.Request) -> httpx.Response:
assert request.url.path.endswith("/getUpdates")
return httpx.Response(200, json={"ok": True, "result": [{"update_id": 1}]})
_install_transport(handler)
client = TelegramClient(token="fake-token")
updates = await client.get_updates(offset=1)
assert updates == [{"update_id": 1}]
async def test_never_logs_or_leaks_token_in_request_url_host() -> None:
"""Sanity: token lives only in the path, base host stays api.telegram.org."""
captured: dict[str, str] = {}
def handler(request: httpx.Request) -> httpx.Response:
captured["url"] = str(request.url)
return httpx.Response(200, json={"ok": True, "result": {}})
_install_transport(handler)
client = TelegramClient(token="super-secret-token")
await client.send_message(chat_id=1, text="hi")
assert "bot" + "super-secret-token" in captured["url"] # goes over the wire, not logged
async def test_copy_message_retries_on_429_then_succeeds() -> None:
calls = {"n": 0}
def handler(request: httpx.Request) -> httpx.Response:
calls["n"] += 1
if calls["n"] == 1:
return httpx.Response(
429,
json={
"ok": False,
"error_code": 429,
"description": "Too Many Requests",
"parameters": {"retry_after": 3},
},
)
return httpx.Response(200, json={"ok": True, "result": {"message_id": 5}})
_install_transport(handler)
client = TelegramClient(token="fake-token")
result = await client.copy_message(chat_id=1, from_chat_id=2, message_id=3)
assert result == {"message_id": 5}
assert calls["n"] == 2
async def test_send_message_retries_on_5xx_then_succeeds() -> None:
calls = {"n": 0}
def handler(request: httpx.Request) -> httpx.Response:
calls["n"] += 1
if calls["n"] < 3:
return httpx.Response(
502, json={"ok": False, "error_code": 502, "description": "bad gw"}
)
return httpx.Response(200, json={"ok": True, "result": {"message_id": 9}})
_install_transport(handler)
client = TelegramClient(token="fake-token")
result = await client.send_message(chat_id=1, text="retrying")
assert result == {"message_id": 9}
assert calls["n"] == 3
async def test_send_message_raises_immediately_on_non_retryable_4xx() -> None:
calls = {"n": 0}
def handler(request: httpx.Request) -> httpx.Response:
calls["n"] += 1
return httpx.Response(
403, json={"ok": False, "error_code": 403, "description": "Forbidden: bot blocked"}
)
_install_transport(handler)
client = TelegramClient(token="fake-token")
with pytest.raises(TelegramApiError) as exc_info:
await client.send_message(chat_id=1, text="hi")
assert exc_info.value.error_code == 403
assert calls["n"] == 1 # НЕ ретраится
async def test_copy_message_gives_up_after_max_retries_on_persistent_5xx() -> None:
def handler(request: httpx.Request) -> httpx.Response:
return httpx.Response(500, json={"ok": False, "error_code": 500, "description": "boom"})
_install_transport(handler)
client = TelegramClient(token="fake-token")
with pytest.raises(TelegramApiError) as exc_info:
await client.copy_message(chat_id=1, from_chat_id=2, message_id=3)
assert exc_info.value.error_code == 500
async def test_get_updates_returns_empty_list_on_malformed_result() -> None:
def handler(request: httpx.Request) -> httpx.Response:
return httpx.Response(200, json={"ok": True, "result": "not-a-list"})
_install_transport(handler)
client = TelegramClient(token="fake-token")
assert await client.get_updates(offset=1) == []
async def test_optional_thread_and_reply_params_omitted_when_falsy() -> None:
captured: dict[str, Any] = {}
def handler(request: httpx.Request) -> httpx.Response:
import json as _json
captured["body"] = _json.loads(request.content.decode("utf-8"))
return httpx.Response(200, json={"ok": True, "result": {"message_id": 1}})
_install_transport(handler)
client = TelegramClient(token="fake-token")
await client.copy_message(chat_id=1, from_chat_id=2, message_id=3, message_thread_id=0)
assert "message_thread_id" not in captured["body"]
async def test_network_error_log_names_the_exception_type(caplog) -> None:
"""В логе сетевого сбоя обязан быть ТИП исключения, а не только текст (#3156).
У `httpx.ReadError` и `httpx.ConnectError` текст обычно пуст, и строка
вырождалась в «network error (попытка 1/3): — retry через 2s»: после
двоеточия пустота. По ней невозможно отличить таймаут от обрыва соединения,
то есть 23 срабатывания в сутки на проде не давали ни одной зацепки.
Проверяем именно пустой текст — на непустом дефект и не проявлялся.
"""
calls = {"n": 0}
def handler(request: httpx.Request) -> httpx.Response:
calls["n"] += 1
if calls["n"] == 1:
raise httpx.ReadError("")
return httpx.Response(200, json={"ok": True, "result": []})
_install_transport(handler)
with caplog.at_level("WARNING", logger="app.services.tgbot.client"):
await TelegramClient(token="t").get_updates(offset=0)
warnings = [r.getMessage() for r in caplog.records if r.levelname == "WARNING"]
assert warnings, "не было предупреждения о сетевом сбое"
assert "ReadError" in warnings[0], (
f"тип исключения не попал в лог, диагностировать нечем: {warnings[0]!r}"
)
assert "network error (попытка 1/" in warnings[0], "формат строки изменился незаметно"
async def test_network_error_log_keeps_text_when_exception_has_one(caplog) -> None:
"""Когда текст у исключения есть — он остаётся, а тип добавляется к нему.
Обратный конец: правка не должна была ЗАМЕНИТЬ текст типом, иначе на
исключениях с внятным сообщением диагностика стала бы беднее прежней.
"""
calls = {"n": 0}
def handler(request: httpx.Request) -> httpx.Response:
calls["n"] += 1
if calls["n"] == 1:
raise httpx.ConnectTimeout("таймаут соединения")
return httpx.Response(200, json={"ok": True, "result": []})
_install_transport(handler)
with caplog.at_level("WARNING", logger="app.services.tgbot.client"):
await TelegramClient(token="t").get_updates(offset=0)
warnings = [r.getMessage() for r in caplog.records if r.levelname == "WARNING"]
assert warnings, "не было предупреждения о сетевом сбое"
assert "ConnectTimeout" in warnings[0], f"нет типа: {warnings[0]!r}"
assert "таймаут соединения" in warnings[0], f"текст исключения потерян: {warnings[0]!r}"
@pytest.mark.parametrize(
("exc_type", "expected_attempts"),
[
(httpx.ConnectTimeout, 2),
(httpx.RemoteProtocolError, 2),
(httpx.ProxyError, 2),
(httpx.DecodingError, 1),
],
)
async def test_request_failure_raises_own_type_not_raw_httpx(
exc_type: type[Exception], expected_attempts: int
) -> None:
"""Любой отказ запроса — наружу свой тип, а не сырой httpx.
Сырой httpx пролетал мимо `except TelegramApiError` во всех трёх HTTP-ручках
и превращался в 500 вместо задуманного 502 (#3456). Первый заход закрыл
только `(TimeoutException, NetworkError)`, а `RemoteProtocolError` («Server
disconnected without sending a response» — бытовой ответ api.telegram.org из
РФ), `ProxyError` и `DecodingError` — сёстры по `TransportError`/
`RequestError`, не наследники `NetworkError`, и дыра оставалась открытой.
Тип отказа при этом терять нельзя — он остаётся в `__cause__`, иначе в
GlitchTip не отличить таймаут соединения от сброса TLS.
"""
def handler(request: httpx.Request) -> httpx.Response:
raise exc_type("сбой транспорта")
_install_transport(handler)
with pytest.raises(TelegramNetworkError) as caught:
await TelegramClient(token="t").send_message(chat_id=-1, text="x", max_retries=1)
assert caught.value.method == "sendMessage"
assert caught.value.attempts == expected_attempts, "число попыток должно попасть в исключение"
assert exc_type.__name__ in caught.value.reason
assert isinstance(caught.value.__cause__, exc_type), "причина потеряна"
async def test_remote_protocol_error_is_retried_but_decoding_error_is_not() -> None:
"""Разница бюджета между двумя `except`: что чинится повтором, а что нет.
`RemoteProtocolError` — «площадка не ответила», ровно как таймаут: повтор
осмыслен, и он наследует уже принятый здесь риск at-least-once (запрос мог
дойти до Telegram, потерялся ответ) — тот же, что у `ReadTimeout`.
`DecodingError` — испорченный ответ / кривая конфигурация: пять попыток с
backoff подвесили бы интерактивную ручку почти на минуту без единого шанса
на успех.
"""
counts: dict[str, int] = {}
async def _attempts_for(exc_type: type[Exception]) -> int:
counts[exc_type.__name__] = 0
def handler(request: httpx.Request) -> httpx.Response:
counts[exc_type.__name__] += 1
raise exc_type("сбой транспорта")
_install_transport(handler)
with pytest.raises(TelegramNetworkError):
await TelegramClient(token="t").send_message(chat_id=-1, text="x", max_retries=2)
return counts[exc_type.__name__]
assert await _attempts_for(httpx.RemoteProtocolError) == 3, (
"обрыв протокола обязан ретраиться наравне с таймаутом"
)
assert await _attempts_for(httpx.DecodingError) == 1, (
"битый ответ ретраить нельзя — повтор не лечит, а бюджет ручки съедает"
)
async def test_network_error_is_not_api_error() -> None:
"""`bridge` разбирает `error_code` (403 «бот заблокирован») — недоступность
площадки в этот разбор попадать не должна, у неё кода ответа нет вовсе."""
def handler(request: httpx.Request) -> httpx.Response:
raise httpx.ConnectTimeout("таймаут соединения")
_install_transport(handler)
with pytest.raises(TelegramNetworkError) as caught:
await TelegramClient(token="t").send_message(chat_id=-1, text="x", max_retries=0)
assert not isinstance(caught.value, TelegramApiError)
# --- раздельные таймауты (#tg-connection-resilience) --------------------------
async def _capture_request_timeout(call) -> dict[str, float]:
captured: dict[str, Any] = {}
def handler(request: httpx.Request) -> httpx.Response:
captured["timeout"] = request.extensions.get("timeout")
return httpx.Response(200, json={"ok": True, "result": []})
_install_transport(handler)
await call(TelegramClient(token="t"))
assert captured["timeout"] is not None, "таймаут не доехал до запроса"
return captured["timeout"]
async def test_long_poll_timeout_applies_to_read_only_not_to_connect() -> None:
"""40 секунд запаса long-poll'а — это бюджет ОТВЕТА, а не установки соединения.
Скаляр в httpx разворачивается в connect=read=write=pool, поэтому
`getUpdates` ждал 40с и коннекта тоже. Живой connect до api.telegram.org из
прод-контейнера — 0.036с; худший цикл из-за этого растягивался на ~174с
(4 попытки × 40с + backoff), и всё это время бот не видел ответов оператора.
"""
timeout = await _capture_request_timeout(lambda tg: tg.get_updates(offset=0, timeout=30))
assert timeout["read"] == 40.0, "long-poll обязан сохранить свои 30+10с на ответ"
assert timeout["connect"] == _CONNECT_TIMEOUT_S, "connect не должен наследовать long-poll"
assert timeout["write"] == _WRITE_TIMEOUT_S
assert timeout["pool"] == _POOL_TIMEOUT_S
async def test_interactive_call_keeps_its_own_narrow_read_budget() -> None:
"""Узкий интерактивный бюджет ручки — тоже read, и он не подменяется дефолтом."""
timeout = await _capture_request_timeout(
lambda tg: tg.send_message(chat_id=1, text="x", timeout=6.0)
)
assert timeout["read"] == 6.0
assert timeout["connect"] == _CONNECT_TIMEOUT_S
# --- переиспользование соединения --------------------------------------------
async def test_two_calls_share_one_httpx_client_and_aclose_closes_it() -> None:
"""Один `httpx.AsyncClient` на жизнь `TelegramClient`, а не на запрос.
Раньше клиент создавался ВНУТРИ цикла ретраев: нулевой keep-alive, полный
TCP+TLS-хендшейк на каждый запрос и на каждую попытку.
"""
def handler(request: httpx.Request) -> httpx.Response:
return httpx.Response(200, json={"ok": True, "result": {"message_id": 1}})
created = _install_transport(handler)
tg = TelegramClient(token="t")
await tg.send_message(chat_id=1, text="раз")
await tg.send_message(chat_id=1, text="два")
assert len(created) == 1, f"клиент пересоздаётся на запрос: {len(created)} штук"
assert not created[0].is_closed
await tg.aclose()
assert created[0].is_closed, "aclose() обязан закрыть пул соединений"
async def test_retries_reuse_the_same_http_client() -> None:
"""Ретраи не пересоздают клиент — иначе повтор платит за хендшейк заново."""
calls = {"n": 0}
def handler(request: httpx.Request) -> httpx.Response:
calls["n"] += 1
if calls["n"] < 3:
return httpx.Response(502, json={"ok": False, "error_code": 502, "description": "gw"})
return httpx.Response(200, json={"ok": True, "result": {"message_id": 7}})
created = _install_transport(handler)
await TelegramClient(token="t").send_message(chat_id=1, text="x")
assert calls["n"] == 3, "бюджет ретраев изменился незаметно"
assert len(created) == 1, "на каждую попытку создаётся новый клиент"
async def test_async_context_manager_closes_client_on_exit() -> None:
def handler(request: httpx.Request) -> httpx.Response:
return httpx.Response(200, json={"ok": True, "result": {"message_id": 1}})
created = _install_transport(handler)
async with TelegramClient(token="t") as tg:
await tg.send_message(chat_id=1, text="x")
assert len(created) == 1
assert created[0].is_closed