diff --git a/backend/app/api/v1/admin_scrape.py b/backend/app/api/v1/admin_scrape.py index 713a3e82..ffa30d0c 100644 --- a/backend/app/api/v1/admin_scrape.py +++ b/backend/app/api/v1/admin_scrape.py @@ -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, diff --git a/backend/tests/test_2464c_timeout_guard_returns.py b/backend/tests/test_2464c_timeout_guard_returns.py index 867e5b8a..6ca04c9a 100644 --- a/backend/tests/test_2464c_timeout_guard_returns.py +++ b/backend/tests/test_2464c_timeout_guard_returns.py @@ -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