Замер 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
486 lines
21 KiB
Python
486 lines
21 KiB
Python
"""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"]
|