"""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"]