fix(ptica): проба глубины очереди перестаёт висеть без таймаута (#2464) #2941
2 changed files with 99 additions and 13 deletions
|
|
@ -191,6 +191,7 @@ def queue_status(
|
||||||
return None
|
return None
|
||||||
|
|
||||||
deadline = time.monotonic() + 0.8
|
deadline = time.monotonic() + 0.8
|
||||||
|
|
||||||
# #2464-C: НЕ `with ThreadPoolExecutor(...)`. Его __exit__ зовёт
|
# #2464-C: НЕ `with ThreadPoolExecutor(...)`. Его __exit__ зовёт
|
||||||
# shutdown(wait=True), поэтому обещанные ~600 мс худшего случая не выполнялись:
|
# shutdown(wait=True), поэтому обещанные ~600 мс худшего случая не выполнялись:
|
||||||
# result(timeout=...) переставал ждать значение, а выход из блока всё равно
|
# result(timeout=...) переставал ждать значение, а выход из блока всё равно
|
||||||
|
|
@ -200,10 +201,25 @@ def queue_status(
|
||||||
# ЧЕСТНАЯ ЦЕНА: shutdown(wait=False) оставляет зависший поток дорабатывать в
|
# ЧЕСТНАЯ ЦЕНА: shutdown(wait=False) оставляет зависший поток дорабатывать в
|
||||||
# фоне. Ограничиваем ЗАПРОС, не процесс — потоки пула не-демоны и джойнятся в
|
# фоне. Ограничиваем ЗАПРОС, не процесс — потоки пула не-демоны и джойнятся в
|
||||||
# atexit. Размен осознанный: висящий поллинг-эндпоинт хуже висящего потока.
|
# 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:
|
try:
|
||||||
f_reserved = ex.submit(_safe, inspect.reserved)
|
f_reserved = ex.submit(_safe, inspect.reserved)
|
||||||
f_ping = ex.submit(_safe, inspect.ping)
|
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:
|
try:
|
||||||
reserved_raw = f_reserved.result(timeout=max(0.1, deadline - time.monotonic()))
|
reserved_raw = f_reserved.result(timeout=max(0.1, deadline - time.monotonic()))
|
||||||
except concurrent.futures.TimeoutError:
|
except concurrent.futures.TimeoutError:
|
||||||
|
|
@ -212,24 +228,20 @@ def queue_status(
|
||||||
ping_resp = f_ping.result(timeout=max(0.1, deadline - time.monotonic()))
|
ping_resp = f_ping.result(timeout=max(0.1, deadline - time.monotonic()))
|
||||||
except concurrent.futures.TimeoutError:
|
except concurrent.futures.TimeoutError:
|
||||||
ping_resp = None
|
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:
|
finally:
|
||||||
ex.shutdown(wait=False, cancel_futures=True)
|
ex.shutdown(wait=False, cancel_futures=True)
|
||||||
|
|
||||||
reserved = _flatten(reserved_raw)
|
reserved = _flatten(reserved_raw)
|
||||||
workers = list((ping_resp or {}).keys())
|
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 {
|
return {
|
||||||
"workers": workers,
|
"workers": workers,
|
||||||
"queue_depth": queue_depth,
|
"queue_depth": queue_depth,
|
||||||
|
|
|
||||||
|
|
@ -156,3 +156,77 @@ def test_queue_status_returns_when_broker_hangs(monkeypatch) -> None:
|
||||||
f"ручка вернулась за {elapsed:.2f} с при обещанных ~0.6 с — "
|
f"ручка вернулась за {elapsed:.2f} с при обещанных ~0.6 с — "
|
||||||
"значит выход из блока ждал зависший inspect, и обещание докстроки ложно"
|
"значит выход из блока ждал зависший 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