fix(nspd_geo): split zombie thresholds running=6m / paused=30m (#1215)
All checks were successful
CI / changes (push) Successful in 6s
CI / frontend-tests (push) Has been skipped
CI / changes (pull_request) Successful in 5s
CI / frontend-tests (pull_request) Has been skipped
CI / backend-tests (push) Successful in 6m24s
CI / backend-tests (pull_request) Successful in 6m23s
All checks were successful
CI / changes (push) Successful in 6s
CI / frontend-tests (push) Has been skipped
CI / changes (pull_request) Successful in 5s
CI / frontend-tests (pull_request) Has been skipped
CI / backend-tests (push) Successful in 6m24s
CI / backend-tests (pull_request) Successful in 6m23s
В WAF-ветке `nspd_geo` worker спит до 240с (4 мин при consecutive_waf=8),
heartbeat коммитится ДО сна. Старый `cleanup_zombies` (`* * * * *`) с
порогом `INTERVAL '2 minutes'` для status IN ('running', 'paused')
гарантированно ре-enqueue'ил **живой** WAF-job:
- Два worker'а параллельно долбили забаненный NSPD-сервис, углубляя бан.
- SELECT pending без `FOR UPDATE SKIP LOCKED` → дубли запросов.
- Counter'ы done/failed двух инстансов затираются друг другом.
Бонус-баг: status='paused' (после 8 WAF подряд) воскресал через минуту,
аннулируя WAF-защиту.
Patch: разделил пороги через module-level константы и параметризованный
SQL (CAST(:x AS interval) per psycopg v3 rule):
- `_ZOMBIE_RUNNING_THRESHOLD = "6 minutes"` — выше max WAF backoff (4 мин).
- `_ZOMBIE_PAUSED_THRESHOLD = "30 minutes"` — реальная WAF-пауза.
Beat-расписание (1 мин tick) не тронуто. WAF-логика / heartbeat-in-loop
(вариант B) / FOR UPDATE SKIP LOCKED (вариант «идеальный») вне scope.
4 новых юнит-теста (test_nspd_geo.py): psycopg v3 guard, оба порога
биндятся, два предиката вместо одного, paused > running.
18/18 nspd_geo тестов зелёные. ruff clean.
Closes #1215
This commit is contained in:
parent
323e399593
commit
787f8a25b9
2 changed files with 155 additions and 9 deletions
|
|
@ -42,6 +42,15 @@ HEARTBEAT_EVERY = 5
|
||||||
# WAF backoff (seconds) после первого 403/429
|
# WAF backoff (seconds) после первого 403/429
|
||||||
WAF_BACKOFF_BASE_S = 30
|
WAF_BACKOFF_BASE_S = 30
|
||||||
|
|
||||||
|
# Zombie-thresholds для cleanup_zombies (issue #1215).
|
||||||
|
# WAF-ветка процесса спит до WAF_BACKOFF_BASE_S * 8 = 240s (4 мин) внутри одной итерации —
|
||||||
|
# heartbeat за это время НЕ обновляется. Старый порог 2 мин ре-enqueue'ил живые WAF-jobs →
|
||||||
|
# два worker'а параллельно долбили забаненный NSPD-сервис, углубляя бан + затирали счётчики.
|
||||||
|
# Поднимаем running-порог выше max WAF-backoff (6 мин > 4 мин). Для 'paused' (consecutive_waf>=8)
|
||||||
|
# нужен большой cooldown, иначе минута beat сразу аннулирует WAF-паузу.
|
||||||
|
_ZOMBIE_RUNNING_THRESHOLD = "6 minutes"
|
||||||
|
_ZOMBIE_PAUSED_THRESHOLD = "30 minutes"
|
||||||
|
|
||||||
|
|
||||||
# ── Helpers для job/target lifecycle ────────────────────────────────────────
|
# ── Helpers для job/target lifecycle ────────────────────────────────────────
|
||||||
|
|
||||||
|
|
@ -562,12 +571,15 @@ def process_nspd_geo_job(self: Any, job_id: int) -> dict[str, Any]:
|
||||||
def cleanup_zombies() -> dict[str, Any]:
|
def cleanup_zombies() -> dict[str, Any]:
|
||||||
"""Periodic zombie cleanup — runs every minute via beat schedule.
|
"""Periodic zombie cleanup — runs every minute via beat schedule.
|
||||||
|
|
||||||
Catches nspd_geo_jobs in status='running' / 'paused' with stale heartbeat
|
Re-enqueues stale nspd_geo_jobs split by status (issue #1215):
|
||||||
(>2 min) and re-enqueues them. Replaces the unreliable worker_ready signal
|
- status='running' с heartbeat старше _ZOMBIE_RUNNING_THRESHOLD (6 мин) →
|
||||||
handler — beat is independent of worker lifecycle and fires every minute.
|
выше max WAF-backoff (240s), не ре-enqueue'ит живые WAF-паузы внутри цикла.
|
||||||
|
- status='paused' с heartbeat старше _ZOMBIE_PAUSED_THRESHOLD (30 мин) →
|
||||||
|
WAF-пауза (consecutive_waf>=8) живёт минимум 30 мин, иначе минутный beat
|
||||||
|
сразу аннулировал бы защиту и worker снова долбил бы забаненный сервис.
|
||||||
|
|
||||||
Idempotent: if no zombies, does nothing. If a job is genuinely active, its
|
Idempotent: если зомби нет, ничего не делает. Активный job с свежим heartbeat
|
||||||
heartbeat will be fresh and the WHERE clause won't match.
|
не матчит WHERE-clause.
|
||||||
"""
|
"""
|
||||||
db = SessionLocal()
|
db = SessionLocal()
|
||||||
resumed: list[int] = []
|
resumed: list[int] = []
|
||||||
|
|
@ -579,11 +591,23 @@ def cleanup_zombies() -> dict[str, Any]:
|
||||||
UPDATE nspd_geo_jobs
|
UPDATE nspd_geo_jobs
|
||||||
SET status = 'queued',
|
SET status = 'queued',
|
||||||
error = COALESCE(error, 'auto-resume by beat cleanup')
|
error = COALESCE(error, 'auto-resume by beat cleanup')
|
||||||
WHERE status IN ('running', 'paused')
|
WHERE (
|
||||||
AND heartbeat_at < NOW() - INTERVAL '2 minutes'
|
status = 'running'
|
||||||
|
AND heartbeat_at
|
||||||
|
< NOW() - CAST(:running_threshold AS interval)
|
||||||
|
)
|
||||||
|
OR (
|
||||||
|
status = 'paused'
|
||||||
|
AND heartbeat_at
|
||||||
|
< NOW() - CAST(:paused_threshold AS interval)
|
||||||
|
)
|
||||||
RETURNING job_id
|
RETURNING job_id
|
||||||
"""
|
"""
|
||||||
)
|
),
|
||||||
|
{
|
||||||
|
"running_threshold": _ZOMBIE_RUNNING_THRESHOLD,
|
||||||
|
"paused_threshold": _ZOMBIE_PAUSED_THRESHOLD,
|
||||||
|
},
|
||||||
)
|
)
|
||||||
.mappings()
|
.mappings()
|
||||||
.all()
|
.all()
|
||||||
|
|
|
||||||
|
|
@ -14,7 +14,15 @@ from __future__ import annotations
|
||||||
from typing import Any
|
from typing import Any
|
||||||
from unittest.mock import MagicMock
|
from unittest.mock import MagicMock
|
||||||
|
|
||||||
from app.workers.tasks.nspd_geo import _save_building, _save_parcel, _save_quarter
|
from app.workers.tasks import nspd_geo as nspd_geo_mod
|
||||||
|
from app.workers.tasks.nspd_geo import (
|
||||||
|
_ZOMBIE_PAUSED_THRESHOLD,
|
||||||
|
_ZOMBIE_RUNNING_THRESHOLD,
|
||||||
|
_save_building,
|
||||||
|
_save_parcel,
|
||||||
|
_save_quarter,
|
||||||
|
cleanup_zombies,
|
||||||
|
)
|
||||||
|
|
||||||
# ── Фикстуры payload (форма ответа _persist_target: {"data": {"features": [...]}}) ──
|
# ── Фикстуры payload (форма ответа _persist_target: {"data": {"features": [...]}}) ──
|
||||||
|
|
||||||
|
|
@ -292,3 +300,117 @@ def test_save_sql_uses_cast_not_double_colon() -> None:
|
||||||
for sql, _params in captured:
|
for sql, _params in captured:
|
||||||
assert "::jsonb" not in sql, f"SQL содержит ::jsonb — нужен CAST(:x AS jsonb): {sql[:80]}"
|
assert "::jsonb" not in sql, f"SQL содержит ::jsonb — нужен CAST(:x AS jsonb): {sql[:80]}"
|
||||||
assert "CAST(:props AS jsonb)" in sql
|
assert "CAST(:props AS jsonb)" in sql
|
||||||
|
|
||||||
|
|
||||||
|
# ── cleanup_zombies thresholds (issue #1215: WAF backoff vs zombie race) ────
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeRowsResult:
|
||||||
|
"""sqlalchemy-like Result, который запоминает (sql, params) и возвращает 0 rows.
|
||||||
|
|
||||||
|
cleanup_zombies использует cain `.mappings().all()` — мокаем под него.
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self, captured: list[tuple[str, dict[str, Any]]], stmt: Any, params: Any) -> None:
|
||||||
|
captured.append((str(stmt), params or {}))
|
||||||
|
self._stmt = stmt
|
||||||
|
self._params = params
|
||||||
|
|
||||||
|
def mappings(self) -> _FakeRowsResult:
|
||||||
|
return self
|
||||||
|
|
||||||
|
def all(self) -> list[dict[str, Any]]:
|
||||||
|
return []
|
||||||
|
|
||||||
|
|
||||||
|
def _capturing_session_factory() -> tuple[MagicMock, list[tuple[str, dict[str, Any]]]]:
|
||||||
|
"""Fake SessionLocal() — возвращает session, чьи execute-calls фиксируются."""
|
||||||
|
captured: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
session = MagicMock()
|
||||||
|
session.execute = lambda stmt, params=None: _FakeRowsResult(captured, stmt, params)
|
||||||
|
session.commit = MagicMock()
|
||||||
|
session.rollback = MagicMock()
|
||||||
|
session.close = MagicMock()
|
||||||
|
|
||||||
|
factory = MagicMock(return_value=session)
|
||||||
|
return factory, captured
|
||||||
|
|
||||||
|
|
||||||
|
def test_cleanup_zombies_sql_uses_cast_as_interval_not_double_colon(
|
||||||
|
monkeypatch: Any,
|
||||||
|
) -> None:
|
||||||
|
"""psycopg v3: только CAST(:x AS interval), запрещён `:x::interval`."""
|
||||||
|
factory, captured = _capturing_session_factory()
|
||||||
|
monkeypatch.setattr(nspd_geo_mod, "SessionLocal", factory)
|
||||||
|
|
||||||
|
cleanup_zombies()
|
||||||
|
|
||||||
|
assert len(captured) == 1, "cleanup_zombies должен выполнить ровно один UPDATE"
|
||||||
|
sql, _params = captured[0]
|
||||||
|
assert "CAST(:running_threshold AS interval)" in sql
|
||||||
|
assert "CAST(:paused_threshold AS interval)" in sql
|
||||||
|
# SQLAlchemy + psycopg v3 ломаются на `:name::type` (backend.md)
|
||||||
|
assert ":running_threshold::" not in sql
|
||||||
|
assert ":paused_threshold::" not in sql
|
||||||
|
|
||||||
|
|
||||||
|
def test_cleanup_zombies_binds_running_and_paused_thresholds(monkeypatch: Any) -> None:
|
||||||
|
"""Issue #1215: running порог > max WAF backoff (240s = 4 min), paused — большой cooldown.
|
||||||
|
|
||||||
|
Контракт thresholds задокументирован в module-level константах.
|
||||||
|
"""
|
||||||
|
factory, captured = _capturing_session_factory()
|
||||||
|
monkeypatch.setattr(nspd_geo_mod, "SessionLocal", factory)
|
||||||
|
|
||||||
|
cleanup_zombies()
|
||||||
|
|
||||||
|
_sql, params = captured[0]
|
||||||
|
assert params["running_threshold"] == _ZOMBIE_RUNNING_THRESHOLD
|
||||||
|
assert params["paused_threshold"] == _ZOMBIE_PAUSED_THRESHOLD
|
||||||
|
# Контракт: running threshold выше старых 2 мин (которые ре-enqueue'или живые WAF-jobs)
|
||||||
|
# и выше max WAF backoff (WAF_BACKOFF_BASE_S * 8 = 240s = 4 мин).
|
||||||
|
assert _ZOMBIE_RUNNING_THRESHOLD == "6 minutes"
|
||||||
|
# Контракт: paused threshold большой (WAF-пауза должна реально удержаться).
|
||||||
|
assert _ZOMBIE_PAUSED_THRESHOLD == "30 minutes"
|
||||||
|
|
||||||
|
|
||||||
|
def test_cleanup_zombies_splits_running_and_paused_in_where_clause(
|
||||||
|
monkeypatch: Any,
|
||||||
|
) -> None:
|
||||||
|
"""SQL разделяет 'running' и 'paused' с разными порогами — не общий IN-фильтр."""
|
||||||
|
factory, captured = _capturing_session_factory()
|
||||||
|
monkeypatch.setattr(nspd_geo_mod, "SessionLocal", factory)
|
||||||
|
|
||||||
|
cleanup_zombies()
|
||||||
|
|
||||||
|
sql, _params = captured[0]
|
||||||
|
executable = _strip_sql_comments(sql)
|
||||||
|
# Старая ветка с общим INTERVAL '2 minutes' для обоих статусов удалена
|
||||||
|
assert "INTERVAL '2 minutes'" not in executable
|
||||||
|
assert "status IN ('running', 'paused')" not in executable
|
||||||
|
# Новая ветка: status='running' и status='paused' в отдельных предикатах
|
||||||
|
assert "status = 'running'" in executable
|
||||||
|
assert "status = 'paused'" in executable
|
||||||
|
# Каждый из двух предикатов сравнивает heartbeat_at с NOW() - threshold
|
||||||
|
assert executable.count("heartbeat_at") >= 2
|
||||||
|
|
||||||
|
|
||||||
|
def test_cleanup_zombies_paused_threshold_strictly_greater_than_running(
|
||||||
|
monkeypatch: Any,
|
||||||
|
) -> None:
|
||||||
|
"""paused cooldown должен быть СТРОГО больше running, иначе пауза аннулируется минутным beat."""
|
||||||
|
factory, captured = _capturing_session_factory()
|
||||||
|
monkeypatch.setattr(nspd_geo_mod, "SessionLocal", factory)
|
||||||
|
|
||||||
|
cleanup_zombies()
|
||||||
|
|
||||||
|
# 6 vs 30 — числовое сравнение префиксов как guard от случайного "6 minutes" == "30 minutes"
|
||||||
|
running_min = int(_ZOMBIE_RUNNING_THRESHOLD.split()[0])
|
||||||
|
paused_min = int(_ZOMBIE_PAUSED_THRESHOLD.split()[0])
|
||||||
|
assert paused_min > running_min, (
|
||||||
|
"paused threshold должен быть строго больше running, "
|
||||||
|
"иначе WAF-пауза не удержится дольше running-cooldown"
|
||||||
|
)
|
||||||
|
# Sanity: SQL действительно отправлен с двумя разными порогами
|
||||||
|
_sql, params = captured[0]
|
||||||
|
assert params["running_threshold"] != params["paused_threshold"]
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue