fix(ptica): защита таймаутом теперь реально ограничивает время (#2464-C) (#2927)
All checks were successful
Deploy / changes (push) Successful in 8s
Deploy / build-frontend (push) Has been skipped
Deploy / build-worker (push) Successful in 4m39s
Deploy / deploy-status (push) Successful in 1s
Deploy / perimeter-smoke (push) Successful in 9s
Deploy / deploy-caddy (push) Has been skipped
Deploy / build-backend (push) Successful in 2m9s
Deploy / deploy (push) Successful in 1m51s

This commit is contained in:
bot-backend 2026-08-19 10:27:58 +00:00
parent cf48e6d6c8
commit e266f29d65
3 changed files with 190 additions and 3 deletions

View file

@ -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())

View file

@ -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)
# ── Общие хелперы фигуры ───────────────────────────────────────────────────────

View 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, и обещание докстроки ложно"
)