From fb52fc095ca3a45fd5e59da4360c2b0baefc7cd0 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Wed, 19 Aug 2026 15:07:20 +0500 Subject: [PATCH] =?UTF-8?q?fix(ptica):=20=D0=B7=D0=B0=D1=89=D0=B8=D1=82?= =?UTF-8?q?=D0=B0=20=D1=82=D0=B0=D0=B9=D0=BC=D0=B0=D1=83=D1=82=D0=BE=D0=BC?= =?UTF-8?q?=20=D1=82=D0=B5=D0=BF=D0=B5=D1=80=D1=8C=20=D1=80=D0=B5=D0=B0?= =?UTF-8?q?=D0=BB=D1=8C=D0=BD=D0=BE=20=D0=BE=D0=B3=D1=80=D0=B0=D0=BD=D0=B8?= =?UTF-8?q?=D1=87=D0=B8=D0=B2=D0=B0=D0=B5=D1=82=20=D0=B2=D1=80=D0=B5=D0=BC?= =?UTF-8?q?=D1=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `with ThreadPoolExecutor(...) as pool:` на выходе зовёт shutdown(wait=True) — он ждёт, пока рабочий поток закончит сам. Поэтому future.result(timeout=T) ограничивал только момент, когда мы перестаём ждать ЗНАЧЕНИЕ, а функция всё равно не возвращалась, пока висел внешний вызов. Два места, один класс: - report_maps._add_basemap — докстрока обещает «недоступный tile-сервер не подвесит воркер», а генерация отчёта стояла столько, сколько стояло зависание; - admin_scrape.queue_status — докстрока обещает «worst-case latency ≈ 600 ms even if no worker is reachable». Ради этого inspect и уносили в поток; ручку фронт опрашивает по таймеру, и висела она вместе с брокером. В обоих случаях except-ветка существовала и выглядела рабочей — она просто ничего не ограничивала. Поэтому тесты меряют ВРЕМЯ ВОЗВРАТА, а не наличие обработчика. ЧЕСТНАЯ ЦЕНА, названная в коде: shutdown(wait=False) оставляет зависший поток дорабатывать в фоне. Ограничивается ЗАПРОС, но не процесс — потоки пула не-демоны и джойнятся в atexit, так что остановка воркера всё ещё может подождать зависший фетч. Это размен «висит генерация отчёта» → «висит один поток в фоне», а не полное устранение. cancel_futures=True снимает только ещё не начатые задачи: начатый поток Python прервать не умеет. Тесты двусторонние: против main падают обе временные проверки, третья — контроль «успешный путь по-прежнему даёт True» — зелёная с обеих сторон. Зависание в тестах ограничено 3 секундами: тест обязан завершаться и на сломанном коде, иначе красный прогон превращается в зависший. Хунк форматирования — не мой: pre-commit ruff v0.7.4 против 0.15.12 (#2864). Refs #2464 --- backend/app/api/v1/admin_scrape.py | 14 +- backend/app/services/exporters/report_maps.py | 21 ++- .../tests/test_2464c_timeout_guard_returns.py | 158 ++++++++++++++++++ 3 files changed, 190 insertions(+), 3 deletions(-) create mode 100644 backend/tests/test_2464c_timeout_guard_returns.py diff --git a/backend/app/api/v1/admin_scrape.py b/backend/app/api/v1/admin_scrape.py index b905d216..713a3e82 100644 --- a/backend/app/api/v1/admin_scrape.py +++ b/backend/app/api/v1/admin_scrape.py @@ -191,7 +191,17 @@ def queue_status( return None deadline = time.monotonic() + 0.8 - with concurrent.futures.ThreadPoolExecutor(max_workers=2) as ex: + # #2464-C: НЕ `with ThreadPoolExecutor(...)`. Его __exit__ зовёт + # shutdown(wait=True), поэтому обещанные ~600 мс худшего случая не выполнялись: + # result(timeout=...) переставал ждать значение, а выход из блока всё равно + # ждал, пока celery inspect отвиснет сам. Для UI-поллинга это ровно та ручка, + # которая обязана возвращаться быстро при недоступном брокере. + # + # ЧЕСТНАЯ ЦЕНА: shutdown(wait=False) оставляет зависший поток дорабатывать в + # фоне. Ограничиваем ЗАПРОС, не процесс — потоки пула не-демоны и джойнятся в + # atexit. Размен осознанный: висящий поллинг-эндпоинт хуже висящего потока. + ex = concurrent.futures.ThreadPoolExecutor(max_workers=2) + try: f_reserved = ex.submit(_safe, inspect.reserved) f_ping = ex.submit(_safe, inspect.ping) try: @@ -202,6 +212,8 @@ def queue_status( ping_resp = f_ping.result(timeout=max(0.1, deadline - time.monotonic())) except concurrent.futures.TimeoutError: ping_resp = None + finally: + ex.shutdown(wait=False, cancel_futures=True) reserved = _flatten(reserved_raw) workers = list((ping_resp or {}).keys()) diff --git a/backend/app/services/exporters/report_maps.py b/backend/app/services/exporters/report_maps.py index 9e8c2097..44fe62db 100644 --- a/backend/app/services/exporters/report_maps.py +++ b/backend/app/services/exporters/report_maps.py @@ -145,9 +145,22 @@ def _add_basemap(ax: Any) -> bool: def _fetch() -> None: cx.add_basemap(ax, crs=_WEB_MERCATOR, source=cx.providers.OpenStreetMap.Mapnik) + # #2464-C: НЕ `with ThreadPoolExecutor(...)`. Его __exit__ зовёт + # shutdown(wait=True) и ждёт, пока рабочий поток реально закончит — то есть + # result(timeout=...) ограничивал момент, когда мы перестаём ждать ЗНАЧЕНИЕ, + # а функция всё равно не возвращалась, пока висел tile-сервер. Заявленный + # «таймаут N секунд» не выполнялся: экспорт стоял столько, сколько стояло + # зависание. + # + # ЧЕСТНАЯ ЦЕНА: shutdown(wait=False) оставляет зависший поток жить до конца + # его собственного вызова. Это ограничивает ЗАПРОС, но не процесс — + # ThreadPoolExecutor держит потоки не-демонами и джойнит их в atexit, так что + # остановка воркера всё ещё может подождать зависший фетч. Меняем «висит + # генерация отчёта» на «висит один поток в фоне» — это осознанный размен, + # а не полное устранение. + pool = ThreadPoolExecutor(max_workers=1) try: - with ThreadPoolExecutor(max_workers=1) as pool: - pool.submit(_fetch).result(timeout=_BASEMAP_TIMEOUT_S) + pool.submit(_fetch).result(timeout=_BASEMAP_TIMEOUT_S) return True except FuturesTimeoutError: logger.warning( @@ -157,6 +170,10 @@ def _add_basemap(ax: Any) -> bool: except Exception as exc: # тайлы недоступны: graceful fallback на белый фон, не валим экспорт logger.warning("report_maps: OSM basemap недоступен (%s) — fallback белый фон", exc) return False + finally: + # cancel_futures=True снимает ещё не начатые задачи; начатую — не отменит + # (Python не умеет прерывать поток), она просто доработает в фоне. + pool.shutdown(wait=False, cancel_futures=True) # ── Общие хелперы фигуры ─────────────────────────────────────────────────────── diff --git a/backend/tests/test_2464c_timeout_guard_returns.py b/backend/tests/test_2464c_timeout_guard_returns.py new file mode 100644 index 00000000..867e5b8a --- /dev/null +++ b/backend/tests/test_2464c_timeout_guard_returns.py @@ -0,0 +1,158 @@ +"""#2464 кластер C: защита таймаутом обязана ОГРАНИЧИВАТЬ время, а не только ожидание. + +`with ThreadPoolExecutor(...) as pool:` на выходе зовёт `shutdown(wait=True)` — +он ждёт, пока рабочий поток реально закончит. Поэтому `future.result(timeout=T)` +ограничивает только момент, когда мы перестаём ждать ЗНАЧЕНИЕ; сама функция всё +равно не вернётся, пока висящий вызов не отвиснет. + +То есть заявленный «таймаут N секунд» не выполняется: при зависшем tile-сервере +экспорт стоит столько, сколько стоит зависание, а не N. + +Тесты меряют ВРЕМЯ ВОЗВРАТА, а не наличие except-ветки: ветка была и раньше, +она просто ничего не ограничивала. +""" + +from __future__ import annotations + +import sys +import threading +import time +import types + +import pytest + + +@pytest.fixture +def _fake_contextily(monkeypatch): + """Подменяет contextily модулем, чей add_basemap ВИСНЕТ. + + Висим ограниченно (3 с), а не вечно: тест обязан завершаться и на сломанном + коде — иначе красный прогон превращается в зависший. + """ + released = threading.Event() + + def _hang(*_a, **_kw): + released.wait(timeout=3.0) + + fake = types.ModuleType("contextily") + fake.add_basemap = _hang + fake.providers = types.SimpleNamespace(OpenStreetMap=types.SimpleNamespace(Mapnik=object())) + monkeypatch.setitem(sys.modules, "contextily", fake) + yield released + released.set() + + +def test_basemap_returns_within_its_own_timeout(_fake_contextily) -> None: + """_add_basemap возвращается около своего таймаута, а не ждёт зависший фетч. + + На main: `with ThreadPoolExecutor` держит выход до конца _hang → ~3 с. + После правки: ~_BASEMAP_TIMEOUT_S. + """ + from app.services.exporters import report_maps + + monkey_timeout = 0.3 + orig = report_maps._BASEMAP_TIMEOUT_S + report_maps._BASEMAP_TIMEOUT_S = monkey_timeout + try: + t0 = time.monotonic() + ok = report_maps._add_basemap(ax=object()) + elapsed = time.monotonic() - t0 + finally: + report_maps._BASEMAP_TIMEOUT_S = orig + + assert ok is False, "зависший фетч не должен считаться успехом" + assert elapsed < 1.5, ( + f"вернулись за {elapsed:.2f} с при таймауте {monkey_timeout} с — " + "значит ждали зависший поток, и заявленный таймаут ничего не ограничивает" + ) + + +def test_basemap_success_path_still_works(monkeypatch) -> None: + """Контроль: рабочий тайл-сервер по-прежнему даёт True. + + Зелёный с обеих сторон правки — доказывает, что правка не превратила + успешный путь в отказ. + """ + calls: list[int] = [] + + fake = types.ModuleType("contextily") + fake.add_basemap = lambda *_a, **_kw: calls.append(1) + fake.providers = types.SimpleNamespace(OpenStreetMap=types.SimpleNamespace(Mapnik=object())) + monkeypatch.setitem(sys.modules, "contextily", fake) + + from app.services.exporters import report_maps + + assert report_maps._add_basemap(ax=object()) is True + assert calls == [1] + + +# ── queue_status: тот же дефект на UI-поллинге ─────────────────────────────── + + +def test_queue_status_returns_when_broker_hangs(monkeypatch) -> None: + """`GET /queue` возвращается по своему дедлайну даже при висящем брокере. + + Докстрока обещает «worst-case latency ≈ 600 ms even if no worker is + reachable» — ради этого inspect и уносили в поток. Но `with + ThreadPoolExecutor(...)` на выходе ждал зависший вызов, и обещание не + выполнялось: ручка, которую фронт опрашивает по таймеру, висела столько, + сколько висел брокер. + """ + import threading + import time + import types + from unittest.mock import MagicMock + + from app.api.v1 import admin_scrape + + released = threading.Event() + + def _hang(*_a, **_kw): + released.wait(timeout=3.0) + return None + + fake_inspect = types.SimpleNamespace(reserved=_hang, ping=_hang) + fake_control = types.SimpleNamespace(inspect=lambda **_kw: fake_inspect) + + # connection_or_acquire — контекст-менеджер; отдаём канал, чей llen мгновенен. + class _Chan: + client = types.SimpleNamespace(llen=lambda _q: 0) + + def __enter__(self): + return self + + def __exit__(self, *a): + return False + + class _Conn: + def channel(self): + return _Chan() + + def __enter__(self): + return self + + def __exit__(self, *a): + return False + + fake_celery = types.SimpleNamespace( + control=fake_control, + connection_or_acquire=lambda: _Conn(), + ) + fake_mod = types.ModuleType("app.workers.celery_app") + fake_mod.celery_app = fake_celery + monkeypatch.setitem(sys.modules, "app.workers.celery_app", fake_mod) + + db = MagicMock() + db.execute.return_value.mappings.return_value.all.return_value = [] + + try: + t0 = time.monotonic() + admin_scrape.queue_status(db=db) + elapsed = time.monotonic() - t0 + finally: + released.set() + + assert elapsed < 2.0, ( + f"ручка вернулась за {elapsed:.2f} с при обещанных ~0.6 с — " + "значит выход из блока ждал зависший inspect, и обещание докстроки ложно" + )