fix(workers): re-raise SoftTimeLimitExceeded в nspd_geo, set paused (#1235)
All checks were successful
CI / changes (pull_request) Successful in 7s
CI / changes (push) Successful in 8s
CI / frontend-tests (pull_request) Has been skipped
CI / frontend-tests (push) Has been skipped
CI / backend-tests (pull_request) Successful in 6m38s
CI / backend-tests (push) Successful in 6m38s
All checks were successful
CI / changes (pull_request) Successful in 7s
CI / changes (push) Successful in 8s
CI / frontend-tests (pull_request) Has been skipped
CI / frontend-tests (push) Has been skipped
CI / backend-tests (pull_request) Successful in 6m38s
CI / backend-tests (push) Successful in 6m38s
This commit is contained in:
parent
0c62d0b0c3
commit
582eb4e8b0
2 changed files with 233 additions and 0 deletions
|
|
@ -23,6 +23,7 @@ import json
|
||||||
import logging
|
import logging
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
|
from celery.exceptions import SoftTimeLimitExceeded
|
||||||
from sqlalchemy import text
|
from sqlalchemy import text
|
||||||
from sqlalchemy.orm import Session
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
|
|
@ -492,6 +493,52 @@ def process_nspd_geo_job(self: Any, job_id: int) -> dict[str, Any]:
|
||||||
return {"job_id": job_id, "paused": True, "reason": "waf"}
|
return {"job_id": job_id, "paused": True, "reason": "waf"}
|
||||||
continue
|
continue
|
||||||
|
|
||||||
|
except SoftTimeLimitExceeded:
|
||||||
|
# Celery soft_time_limit=6h истёк — graceful shutdown ДО generic except.
|
||||||
|
# Issue #1235: `except (NspdLiteError, Exception)` ниже проглатывал бы это
|
||||||
|
# исключение (SoftTimeLimitExceeded унаследован от Exception) → таска
|
||||||
|
# продолжала бы цикл неограниченно дольше 6ч, занимая слот geo-очереди.
|
||||||
|
# Контракт: коммитим финальный heartbeat + статусы targets, помечаем job
|
||||||
|
# 'paused' (auto-resume через cleanup_zombies/worker_ready), re-raise чтобы
|
||||||
|
# Celery освободил воркер и зарегистрировал лимит.
|
||||||
|
logger.warning(
|
||||||
|
"process_nspd_geo_job soft time limit exceeded: job=%s done=%d failed=%d",
|
||||||
|
job_id,
|
||||||
|
done,
|
||||||
|
failed,
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
_heartbeat(
|
||||||
|
db,
|
||||||
|
job_id,
|
||||||
|
targets_done=done,
|
||||||
|
targets_failed=failed,
|
||||||
|
targets_skipped=skipped,
|
||||||
|
requests_count=n_requests,
|
||||||
|
waf_blocked_count=n_waf,
|
||||||
|
)
|
||||||
|
_finish_job(
|
||||||
|
db,
|
||||||
|
job_id,
|
||||||
|
"paused",
|
||||||
|
error=f"soft_time_limit exceeded after {done + failed} targets",
|
||||||
|
)
|
||||||
|
_log(
|
||||||
|
db,
|
||||||
|
job_id,
|
||||||
|
"warn",
|
||||||
|
"paused",
|
||||||
|
f"job paused: Celery soft_time_limit (6h) reached "
|
||||||
|
f"after {done} done / {failed} failed",
|
||||||
|
)
|
||||||
|
except Exception as flush_err:
|
||||||
|
# Финальный flush не должен скрыть SoftTimeLimitExceeded — лог + re-raise.
|
||||||
|
logger.exception(
|
||||||
|
"process_nspd_geo_job: failed to flush state on soft-limit: %s",
|
||||||
|
flush_err,
|
||||||
|
)
|
||||||
|
raise
|
||||||
|
|
||||||
except (NspdLiteError, Exception) as e:
|
except (NspdLiteError, Exception) as e:
|
||||||
logger.warning("fetch failed for %s: %s", cad, e)
|
logger.warning("fetch failed for %s: %s", cad, e)
|
||||||
_log(db, job_id, "warn", "fetch_fail", str(e)[:300], cad_num=cad)
|
_log(db, job_id, "warn", "fetch_fail", str(e)[:300], cad_num=cad)
|
||||||
|
|
@ -556,6 +603,11 @@ def process_nspd_geo_job(self: Any, job_id: int) -> dict[str, Any]:
|
||||||
"requests": n_requests,
|
"requests": n_requests,
|
||||||
"waf_blocked": n_waf,
|
"waf_blocked": n_waf,
|
||||||
}
|
}
|
||||||
|
except SoftTimeLimitExceeded:
|
||||||
|
# Inner handler уже коммитнул heartbeat + статус 'paused' + лог.
|
||||||
|
# Не перетираем 'paused' на 'failed' через outer fallback — просто пробрасываем,
|
||||||
|
# чтобы Celery освободил worker и зарегистрировал срабатывание soft_time_limit.
|
||||||
|
raise
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception("process_nspd_geo_job crashed: job=%s: %s", job_id, e)
|
logger.exception("process_nspd_geo_job crashed: job=%s: %s", job_id, e)
|
||||||
try:
|
try:
|
||||||
|
|
|
||||||
|
|
@ -14,6 +14,9 @@ from __future__ import annotations
|
||||||
from typing import Any
|
from typing import Any
|
||||||
from unittest.mock import MagicMock
|
from unittest.mock import MagicMock
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
from celery.exceptions import SoftTimeLimitExceeded
|
||||||
|
|
||||||
from app.workers.tasks import nspd_geo as nspd_geo_mod
|
from app.workers.tasks import nspd_geo as nspd_geo_mod
|
||||||
from app.workers.tasks.nspd_geo import (
|
from app.workers.tasks.nspd_geo import (
|
||||||
_ZOMBIE_PAUSED_THRESHOLD,
|
_ZOMBIE_PAUSED_THRESHOLD,
|
||||||
|
|
@ -22,6 +25,7 @@ from app.workers.tasks.nspd_geo import (
|
||||||
_save_parcel,
|
_save_parcel,
|
||||||
_save_quarter,
|
_save_quarter,
|
||||||
cleanup_zombies,
|
cleanup_zombies,
|
||||||
|
process_nspd_geo_job,
|
||||||
)
|
)
|
||||||
|
|
||||||
# ── Фикстуры payload (форма ответа _persist_target: {"data": {"features": [...]}}) ──
|
# ── Фикстуры payload (форма ответа _persist_target: {"data": {"features": [...]}}) ──
|
||||||
|
|
@ -414,3 +418,180 @@ def test_cleanup_zombies_paused_threshold_strictly_greater_than_running(
|
||||||
# Sanity: SQL действительно отправлен с двумя разными порогами
|
# Sanity: SQL действительно отправлен с двумя разными порогами
|
||||||
_sql, params = captured[0]
|
_sql, params = captured[0]
|
||||||
assert params["running_threshold"] != params["paused_threshold"]
|
assert params["running_threshold"] != params["paused_threshold"]
|
||||||
|
|
||||||
|
|
||||||
|
# ── SoftTimeLimitExceeded handler (issue #1235) ─────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
class _JobFetchResult:
|
||||||
|
"""sqlalchemy-like Result, возвращающий заданный mapping для SELECT FROM nspd_geo_jobs."""
|
||||||
|
|
||||||
|
def __init__(self, row: dict[str, Any] | None) -> None:
|
||||||
|
self._row = row
|
||||||
|
|
||||||
|
def mappings(self) -> _JobFetchResult:
|
||||||
|
return self
|
||||||
|
|
||||||
|
def first(self) -> dict[str, Any] | None:
|
||||||
|
return self._row
|
||||||
|
|
||||||
|
def scalar_one(self) -> Any:
|
||||||
|
return 1
|
||||||
|
|
||||||
|
|
||||||
|
def _build_soft_limit_session(
|
||||||
|
raise_in: str = "fetch",
|
||||||
|
) -> tuple[MagicMock, list[tuple[str, dict[str, Any]]]]:
|
||||||
|
"""Fake SessionLocal() для теста graceful shutdown по SoftTimeLimitExceeded.
|
||||||
|
|
||||||
|
Контракт: первый SELECT ... FROM nspd_geo_jobs WHERE job_id=... → job row,
|
||||||
|
UPDATE _start_job → OK, INSERT nspd_geo_log → OK, второй SELECT (target loop) →
|
||||||
|
pending target row. Затем fetcher кидает SoftTimeLimitExceeded.
|
||||||
|
|
||||||
|
После raise — handler должен (a) выполнить UPDATE heartbeat + counters,
|
||||||
|
(b) UPDATE nspd_geo_jobs SET status='paused', (c) re-raise.
|
||||||
|
"""
|
||||||
|
captured: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
job_row = {
|
||||||
|
"job_id": 42,
|
||||||
|
"rate_ms": 0, # отключаем sleep между итерациями для теста
|
||||||
|
"targets_done": 0,
|
||||||
|
"targets_failed": 0,
|
||||||
|
"targets_skipped": 0,
|
||||||
|
"requests_count": 0,
|
||||||
|
"waf_blocked_count": 0,
|
||||||
|
}
|
||||||
|
target_row = {
|
||||||
|
"target_id": 7,
|
||||||
|
"cad_num": "66:41:0610029:83",
|
||||||
|
"thematic_id": 1,
|
||||||
|
"attempts": 0,
|
||||||
|
}
|
||||||
|
call_state = {"job_selects": 0, "target_selects": 0}
|
||||||
|
|
||||||
|
def _execute(stmt: Any, params: Any = None) -> Any:
|
||||||
|
sql = str(stmt)
|
||||||
|
captured.append((sql, params or {}))
|
||||||
|
if "FROM nspd_geo_jobs" in sql and "SELECT" in sql:
|
||||||
|
call_state["job_selects"] += 1
|
||||||
|
return _JobFetchResult(job_row)
|
||||||
|
if "FROM nspd_geo_targets" in sql and "SELECT" in sql:
|
||||||
|
call_state["target_selects"] += 1
|
||||||
|
# Первый SELECT — отдаём pending target. Дальше — None (loop break),
|
||||||
|
# но до второй итерации не дойдём — fetcher кинет SoftTimeLimitExceeded.
|
||||||
|
if call_state["target_selects"] == 1:
|
||||||
|
return _JobFetchResult(target_row)
|
||||||
|
return _JobFetchResult(None)
|
||||||
|
if "INSERT INTO nspd_geo_jobs" in sql:
|
||||||
|
return _JobFetchResult({"job_id": 42}) # для scalar_one
|
||||||
|
return MagicMock()
|
||||||
|
|
||||||
|
session = MagicMock()
|
||||||
|
session.execute = _execute
|
||||||
|
session.commit = MagicMock()
|
||||||
|
session.rollback = MagicMock()
|
||||||
|
session.close = MagicMock()
|
||||||
|
return MagicMock(return_value=session), captured
|
||||||
|
|
||||||
|
|
||||||
|
def test_soft_time_limit_exceeded_reraises_not_swallowed(monkeypatch: Any) -> None:
|
||||||
|
"""Issue #1235 root: `except (NspdLiteError, Exception)` проглатывало SoftTimeLimitExceeded.
|
||||||
|
|
||||||
|
После fix — explicit `except SoftTimeLimitExceeded: raise` ДО generic ветки →
|
||||||
|
исключение пробрасывается наружу process_nspd_geo_job, Celery освобождает воркер.
|
||||||
|
"""
|
||||||
|
factory, _captured = _build_soft_limit_session()
|
||||||
|
monkeypatch.setattr(nspd_geo_mod, "SessionLocal", factory)
|
||||||
|
|
||||||
|
def _raise_soft_limit(*_args: Any, **_kwargs: Any) -> Any:
|
||||||
|
raise SoftTimeLimitExceeded()
|
||||||
|
|
||||||
|
monkeypatch.setattr(nspd_geo_mod, "fetch_via_rosreestr2coord", _raise_soft_limit)
|
||||||
|
|
||||||
|
with pytest.raises(SoftTimeLimitExceeded):
|
||||||
|
process_nspd_geo_job(42)
|
||||||
|
|
||||||
|
|
||||||
|
def test_soft_time_limit_exceeded_sets_paused_status(monkeypatch: Any) -> None:
|
||||||
|
"""Handler коммитит status='paused' перед re-raise — иначе job навечно 'running'
|
||||||
|
с stale heartbeat → cleanup_zombies ре-enqueue'ит через running-threshold (6 мин),
|
||||||
|
но мы хотим явный paused + 30 мин cooldown.
|
||||||
|
"""
|
||||||
|
factory, captured = _build_soft_limit_session()
|
||||||
|
monkeypatch.setattr(nspd_geo_mod, "SessionLocal", factory)
|
||||||
|
monkeypatch.setattr(
|
||||||
|
nspd_geo_mod,
|
||||||
|
"fetch_via_rosreestr2coord",
|
||||||
|
MagicMock(side_effect=SoftTimeLimitExceeded()),
|
||||||
|
)
|
||||||
|
|
||||||
|
with pytest.raises(SoftTimeLimitExceeded):
|
||||||
|
process_nspd_geo_job(42)
|
||||||
|
|
||||||
|
# Должен быть UPDATE с status = :st где :st == 'paused'
|
||||||
|
paused_updates = [
|
||||||
|
(sql, params)
|
||||||
|
for sql, params in captured
|
||||||
|
if "UPDATE nspd_geo_jobs" in sql and "status = :st" in sql and params.get("st") == "paused"
|
||||||
|
]
|
||||||
|
assert paused_updates, (
|
||||||
|
"SoftTimeLimitExceeded handler должен выставить status='paused' до re-raise; "
|
||||||
|
f"captured updates: {[(s[:80], p) for s, p in captured if 'UPDATE nspd_geo_jobs' in s]}"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_soft_time_limit_exceeded_flushes_heartbeat(monkeypatch: Any) -> None:
|
||||||
|
"""Handler коммитит финальный heartbeat (counters + timestamp) до raise.
|
||||||
|
|
||||||
|
Иначе после паузы job выглядит в UI с устаревшими счётчиками — admin не видит
|
||||||
|
сколько таргетов реально обработано.
|
||||||
|
"""
|
||||||
|
factory, captured = _build_soft_limit_session()
|
||||||
|
monkeypatch.setattr(nspd_geo_mod, "SessionLocal", factory)
|
||||||
|
monkeypatch.setattr(
|
||||||
|
nspd_geo_mod,
|
||||||
|
"fetch_via_rosreestr2coord",
|
||||||
|
MagicMock(side_effect=SoftTimeLimitExceeded()),
|
||||||
|
)
|
||||||
|
|
||||||
|
with pytest.raises(SoftTimeLimitExceeded):
|
||||||
|
process_nspd_geo_job(42)
|
||||||
|
|
||||||
|
heartbeat_updates = [
|
||||||
|
(sql, params)
|
||||||
|
for sql, params in captured
|
||||||
|
if "heartbeat_at = NOW()" in sql and "targets_done" in sql
|
||||||
|
]
|
||||||
|
assert heartbeat_updates, (
|
||||||
|
"SoftTimeLimitExceeded handler должен flush'нуть heartbeat с counters перед raise"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_soft_time_limit_exceeded_does_not_overwrite_paused_with_failed(
|
||||||
|
monkeypatch: Any,
|
||||||
|
) -> None:
|
||||||
|
"""Outer `except Exception` НЕ должен затереть 'paused' на 'failed' для SoftTimeLimitExceeded.
|
||||||
|
|
||||||
|
Issue #1235 corollary: если inner handler выставил 'paused' и raise — outer
|
||||||
|
fallback не должен мгновенно переписать status='failed' через _finish_job.
|
||||||
|
"""
|
||||||
|
factory, captured = _build_soft_limit_session()
|
||||||
|
monkeypatch.setattr(nspd_geo_mod, "SessionLocal", factory)
|
||||||
|
monkeypatch.setattr(
|
||||||
|
nspd_geo_mod,
|
||||||
|
"fetch_via_rosreestr2coord",
|
||||||
|
MagicMock(side_effect=SoftTimeLimitExceeded()),
|
||||||
|
)
|
||||||
|
|
||||||
|
with pytest.raises(SoftTimeLimitExceeded):
|
||||||
|
process_nspd_geo_job(42)
|
||||||
|
|
||||||
|
failed_updates = [
|
||||||
|
params
|
||||||
|
for sql, params in captured
|
||||||
|
if "UPDATE nspd_geo_jobs" in sql and "status = :st" in sql and params.get("st") == "failed"
|
||||||
|
]
|
||||||
|
assert not failed_updates, (
|
||||||
|
"outer except не должен затирать 'paused' на 'failed' для SoftTimeLimitExceeded; "
|
||||||
|
f"got failed updates: {failed_updates}"
|
||||||
|
)
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue