gendesign/tradein-mvp/backend/tests/services/tgbot/test_client.py
bot-backend e0564d12fe feat(tg): продуктовый Bot API трафик уходит через ретранслятор на Beget
Замер 12.09.2026, оба хоста в одни и те же минуты: getMe из tradein-tgbot
на Selectel — 9 успешных из 12, три ConnectTimeout; TCP-443 до адреса,
резолвящегося на Selectel (149.154.167.220) — 5 из 6; TCP-443 до адреса,
резолвящегося на Beget (149.154.166.110) — 8 из 8. За сутки в логе бота
508 строк network error, за 30 дней 92 обрыва итерации poll loop. Значит:
путь до Telegram с Selectel лоссовый, с Beget чистый — Alertmanager (живёт
на Beget) шлёт в тот же чат без проблем, а бот поддержки на Selectel часть
отправок теряет.

Добавлен ops/metrics/tg-relay — stdlib-only HTTP-сервис (тот же принцип,
что у alert-ack: без зависимостей, поднимается даже когда всё остальное
сломано), проксирует Bot API целиком (метод, путь, тело — sendMessage,
copyMessage, getUpdates) на api.telegram.org. Токен из пути не логируется:
log_request переопределён полностью, путь редактируется до записи в лог.
Аутентификация — общий секрет в X-Relay-Secret, по образцу
X-Internal-Auth-Secret из этого же стека.

Клиент (tgbot/client.py) при транспортном отказе похода на ретранслятор
делает одну попытку напрямую к api.telegram.org — хуже прямого пути быть
не должно ни при каких условиях. Пустой TELEGRAM_RELAY_BASE_URL — прежнее
поведение без изменений, это и есть механизм отката.

Refs #3471
2026-09-12 14:16:05 +03:00

486 lines
21 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
# ── Ретранслятор через Beget: фолбэк на прямой путь (#3471) ────────────────
#
# Единственный `httpx.AsyncClient` (общий для обоих адресов через тот же
# MockTransport) позволяет различать «запрос на ретранслятор» от «запрос
# напрямую» по хосту в `request.url.host` — так проверяется, что при отказе
# ретранслятора клиент реально уходит на api.telegram.org, а не молчит.
async def test_relay_transport_failure_falls_back_to_direct_once() -> None:
"""Отказ ретранслятора не должен ронять бота: одна попытка напрямую."""
seen_hosts: list[str] = []
def handler(request: httpx.Request) -> httpx.Response:
seen_hosts.append(request.url.host)
if request.url.host == "relay.example.com":
raise httpx.ConnectError("connection refused", request=request)
return httpx.Response(200, json={"ok": True, "result": {}})
_install_transport(handler)
client = TelegramClient(
token="fake-token", relay_base_url="https://relay.example.com", relay_secret="shh"
)
result = await client.send_message(chat_id=1, text="hi")
assert result == {}
assert seen_hosts == ["relay.example.com", "api.telegram.org"]
async def test_relay_secret_header_sent_only_to_relay_not_to_direct_fallback() -> None:
captured: list[str | None] = []
def handler(request: httpx.Request) -> httpx.Response:
captured.append(request.headers.get("X-Relay-Secret"))
if request.url.host == "relay.example.com":
raise httpx.ConnectError("connection refused", request=request)
return httpx.Response(200, json={"ok": True, "result": {}})
_install_transport(handler)
client = TelegramClient(
token="fake-token", relay_base_url="https://relay.example.com", relay_secret="shh"
)
await client.send_message(chat_id=1, text="hi")
assert captured == ["shh", None], "секрет ретранслятора не должен уходить напрямую в Telegram"
async def test_no_relay_configured_behaves_exactly_as_direct_path_before() -> None:
"""Пустая настройка — это откат: поведение НЕ должно отличаться от прежнего."""
def handler(request: httpx.Request) -> httpx.Response:
assert request.url.host == "api.telegram.org"
assert "X-Relay-Secret" not in request.headers
return httpx.Response(200, json={"ok": True, "result": {}})
_install_transport(handler)
result = await TelegramClient(token="fake-token").send_message(chat_id=1, text="hi")
assert result == {}
async def test_relay_success_never_touches_direct_host() -> None:
seen_hosts: list[str] = []
def handler(request: httpx.Request) -> httpx.Response:
seen_hosts.append(request.url.host)
return httpx.Response(200, json={"ok": True, "result": {}})
_install_transport(handler)
client = TelegramClient(
token="fake-token", relay_base_url="https://relay.example.com", relay_secret="shh"
)
await client.send_message(chat_id=1, text="hi")
assert seen_hosts == ["relay.example.com"]