fix(ptica): проба глубины очереди перестаёт висеть без таймаута (#2464)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 10s
CI Trade-In / backend-tests (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Successful in 3m0s
CI / backend-tests (pull_request) Successful in 17m37s

Ручка queue_status обещает в докстроке «worst-case latency ≈ 600 ms even if no
worker is reachable» — ради этого celery inspect уносили в поток под дедлайн
0.8 с, и #2927 починил там executor, который сводил защиту на нет.

Но следом шёл третий шаг — проба глубины очереди — СИНХРОННО и без таймаута
вовсе:

    with celery_app.connection_or_acquire() as conn:
        with conn.channel() as channel:
            queue_depth = channel.client.llen("celery")

По висящему сокету (не «connection refused», а чёрная дыра) это не возвращается
никогда. То есть обещание докстроки ломалось на последнем шаге, и ручка, которую
админ-UI опрашивает по таймеру, висела столько, сколько висел брокер.

Проба уходит в тот же пул под тот же дедлайн. Отправляется ДО чтения результатов
inspect, а не после: иначе к моменту её старта бюджет уже израсходован, и ей
достаётся только нижняя граница max(0.1, ...).

Существующий тест на зависание брокера этот случай не покрывал — в нём `llen`
мгновенный, виснут только inspect'ы.

Двусторонняя проверка ПО ВРЕМЕНИ: на origin/main ручка возвращается за 3.01 с
(ровно длительность искусственного зависания), с правкой — меньше 2 с. Тест
дополнительно требует queue_depth is None: не смогли измерить — отдаём None, а не
выдуманный ноль.

pytest -k "admin_scrape or queue or 2464c": 35 passed, 1 skipped, rc=0
pytest tests/api/v1: 354 passed, 1 skipped, rc=0
Тест перепрогнан после правок pre-commit ruff-format.
This commit is contained in:
bot-backend 2026-08-19 21:33:33 +05:00
parent f3626540fc
commit 98e5229277
2 changed files with 99 additions and 13 deletions

View file

@ -191,6 +191,7 @@ def queue_status(
return None
deadline = time.monotonic() + 0.8
# #2464-C: НЕ `with ThreadPoolExecutor(...)`. Его __exit__ зовёт
# shutdown(wait=True), поэтому обещанные ~600 мс худшего случая не выполнялись:
# result(timeout=...) переставал ждать значение, а выход из блока всё равно
@ -200,10 +201,25 @@ def queue_status(
# ЧЕСТНАЯ ЦЕНА: shutdown(wait=False) оставляет зависший поток дорабатывать в
# фоне. Ограничиваем ЗАПРОС, не процесс — потоки пула не-демоны и джойнятся в
# atexit. Размен осознанный: висящий поллинг-эндпоинт хуже висящего потока.
ex = concurrent.futures.ThreadPoolExecutor(max_workers=2)
def _probe_queue_depth() -> int | None:
with celery_app.connection_or_acquire() as conn:
with conn.channel() as channel:
return channel.client.llen("celery")
ex = concurrent.futures.ThreadPoolExecutor(max_workers=3)
try:
f_reserved = ex.submit(_safe, inspect.reserved)
f_ping = ex.submit(_safe, inspect.ping)
# #2464: проба глубины очереди раньше шла СИНХРОННО и без таймаута вовсе —
# `connection_or_acquire()` + `llen` по висящему сокету не возвращаются
# никогда. Дедлайн 0.8 с выше ограничивал только inspect, а ручка всё равно
# висела столько, сколько висел брокер: обещание докстроки не выполнялось на
# последнем шаге. Отправляем в тот же пул под тот же дедлайн.
#
# Отправляем ДО чтения результатов, а не после: иначе к моменту старта пробы
# бюджет уже израсходован inspect'ами и ей досталась бы только нижняя
# граница max(0.1, ...).
f_queue = ex.submit(_safe, _probe_queue_depth)
try:
reserved_raw = f_reserved.result(timeout=max(0.1, deadline - time.monotonic()))
except concurrent.futures.TimeoutError:
@ -212,24 +228,20 @@ def queue_status(
ping_resp = f_ping.result(timeout=max(0.1, deadline - time.monotonic()))
except concurrent.futures.TimeoutError:
ping_resp = None
# 3) Pending in broker queue (not yet picked up by any worker).
# Деградация до None намеренна (см. _safe): недоступный брокер не должен
# ронять UI-поллинг, но обязан быть виден в логах.
try:
queue_depth: int | None = f_queue.result(timeout=max(0.1, deadline - time.monotonic()))
except concurrent.futures.TimeoutError:
logger.warning("queue_status: broker queue_depth probe timed out")
queue_depth = None
finally:
ex.shutdown(wait=False, cancel_futures=True)
reserved = _flatten(reserved_raw)
workers = list((ping_resp or {}).keys())
# 3) Pending in broker queue (not yet picked up by any worker).
queue_depth: int | None = None
try:
with celery_app.connection_or_acquire() as conn:
with conn.channel() as channel:
queue_depth = channel.client.llen("celery")
except Exception:
# Намеренная деградация для UI-poll; логируем чтобы недоступный broker
# не был невидим в логах (см. .claude/rules/backend.md).
logger.warning("queue_status: broker queue_depth probe failed", exc_info=True)
queue_depth = None
return {
"workers": workers,
"queue_depth": queue_depth,

View file

@ -156,3 +156,77 @@ def test_queue_status_returns_when_broker_hangs(monkeypatch) -> None:
f"ручка вернулась за {elapsed:.2f} с при обещанных ~0.6 с"
"значит выход из блока ждал зависший inspect, и обещание докстроки ложно"
)
def test_queue_status_returns_when_the_queue_probe_itself_hangs(monkeypatch) -> None:
"""Виснет НЕ inspect, а сама проба глубины очереди.
Тест выше подменяет `llen` мгновенным, поэтому этот случай не покрывал.
А именно он ломал обещание докстроки на последнем шаге: `inspect` был
ограничен дедлайном 0.8 с, а `connection_or_acquire()` + `llen` шли синхронно
и без таймаута вовсе. По висящему сокету (не «connection refused», а
чёрная дыра) ручка не возвращалась никогда.
Как и выше, висим ограниченно (3 с): тест обязан завершаться и на сломанном
коде, иначе красный прогон превращается в зависший.
"""
import threading
import time
import types
from unittest.mock import MagicMock
from app.api.v1 import admin_scrape
released = threading.Event()
# inspect отвечает мгновенно — весь бюджет остаётся пробе очереди.
fake_inspect = types.SimpleNamespace(reserved=lambda: None, ping=lambda: {})
fake_control = types.SimpleNamespace(inspect=lambda **_kw: fake_inspect)
def _hanging_llen(_queue: str) -> int:
released.wait(timeout=3.0)
return 0
class _Chan:
client = types.SimpleNamespace(llen=_hanging_llen)
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()
out = admin_scrape.queue_status(db=db)
elapsed = time.monotonic() - t0
finally:
released.set()
assert elapsed < 2.0, (
f"ручка вернулась за {elapsed:.2f} с при обещанных ~0.8 с"
"проба глубины очереди по-прежнему идёт без таймаута"
)
# Деградация честная: не смогли измерить — отдаём None, а не выдуманный ноль.
assert out["queue_depth"] is None