fix(ptica): проба глубины очереди перестаёт висеть без таймаута (#2464) #2941
2 changed files with 99 additions and 13 deletions
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue