fix(ptica): защита таймаутом теперь реально ограничивает время (#2464-C) #2927
3 changed files with 190 additions and 3 deletions
|
|
@ -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())
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
||||
|
||||
# ── Общие хелперы фигуры ───────────────────────────────────────────────────────
|
||||
|
|
|
|||
158
backend/tests/test_2464c_timeout_guard_returns.py
Normal file
158
backend/tests/test_2464c_timeout_guard_returns.py
Normal file
|
|
@ -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, и обещание докстроки ложно"
|
||||
)
|
||||
Loading…
Add table
Reference in a new issue