fix(ptica): resume_geo_job больше не возобновляет что попало (#2464)
Some checks failed
CI Trade-In / changes (pull_request) Successful in 7s
CI / changes (pull_request) Successful in 10s
CI Trade-In / backend-tests (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Failing after 2m19s
CI / backend-tests (pull_request) Successful in 17m5s
Some checks failed
CI Trade-In / changes (pull_request) Successful in 7s
CI / changes (pull_request) Successful in 10s
CI Trade-In / backend-tests (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Failing after 2m19s
CI / backend-tests (pull_request) Successful in 17m5s
UPDATE шёл БЕЗ фильтра статуса — в отличие от соседнего cancel_geo_job, который
фильтрует явно (`AND status IN ('queued','running','paused')`):
UPDATE nspd_geo_jobs SET status='queued', error=NULL WHERE job_id=:id
Следствия: завершённую задачу можно было перевести обратно в 'queued' и прогнать
заново, затирая результат; уже бегущую — поставить в очередь второй раз, получив
двух воркеров на один job_id и лишние запросы к НСПД, у которого WAF.
Замер на проде 19.08: все 66 задач в терминальных статусах (61 done, 5 cancelled).
То есть resume на ЛЮБУЮ существующую делал ровно то, чего не должен.
Второе: ручка возвращала resumed=True всегда, независимо от того, изменилось ли
что-нибудь. Теперь ответ отражает факт — не подошёл статус, значит resumed=False,
текущий статус и причина в ответе, задача НЕ ставится в очередь.
'cancelled' оставлен возобновляемым намеренно: cancel — ручное действие оператора,
и без этого отменённая по ошибке задача не восстанавливалась бы никак.
Тесты: 4 красных на origin/main, главный — «AssertionError: UPDATE без фильтра
статуса — возобновляется что угодно».
Контроль пришлось переделать: первая версия проверяла и новый ключ `status`, из-за
чего падала на origin/main с KeyError, то есть по причине «в ответе нет поля», а не
«законный путь сломан». Разделено: контроль смотрит только resumed и зелёный по обе
стороны, новый ключ проверяется отдельным тестом.
pytest tests/api/v1: 359 passed, 1 skipped, rc=0
This commit is contained in:
parent
56868f2bde
commit
af4d2a1853
2 changed files with 157 additions and 5 deletions
|
|
@ -1142,18 +1142,63 @@ def resume_geo_job(
|
|||
job_id: int,
|
||||
db: Annotated[Session, Depends(get_db)],
|
||||
) -> dict[str, Any]:
|
||||
"""Re-enqueue paused/failed job. Resume idempotent через pending targets."""
|
||||
"""Re-enqueue задачу из НЕзавершённого состояния (paused / failed / cancelled).
|
||||
|
||||
#2464: UPDATE шёл БЕЗ фильтра статуса — в отличие от соседнего cancel_geo_job,
|
||||
который фильтрует явно. Из-за этого «возобновить» можно было завершённую задачу
|
||||
(status='done' → снова 'queued' и повторный прогон, затирая результат) и уже
|
||||
бегущую (второй worker на тот же job_id — лишние запросы к НСПД, у которого WAF).
|
||||
|
||||
Замер на проде 19.08: все 66 задач в терминальных статусах — 61 done, 5
|
||||
cancelled. То есть resume на ЛЮБУЮ существующую делал ровно то, чего не должен.
|
||||
|
||||
Второе: ручка возвращала resumed=True всегда, независимо от того, изменилось ли
|
||||
что-нибудь. Теперь ответ отражает факт: не подошёл статус → resumed=False,
|
||||
текущий статус в ответе, задача НЕ ставится в очередь.
|
||||
|
||||
'cancelled' оставлен возобновляемым намеренно: cancel_geo_job — ручное действие
|
||||
оператора, и без этого отменённая по ошибке задача не восстанавливалась бы никак.
|
||||
"""
|
||||
from app.services.job_settings import get_setting_value
|
||||
from app.workers.tasks.nspd_geo import process_nspd_geo_job
|
||||
|
||||
db.execute(
|
||||
text("UPDATE nspd_geo_jobs SET status='queued', error=NULL WHERE job_id=:id"),
|
||||
{"id": job_id},
|
||||
row = (
|
||||
db.execute(
|
||||
text(
|
||||
"""
|
||||
UPDATE nspd_geo_jobs SET status='queued', error=NULL
|
||||
WHERE job_id = :id AND status IN ('paused','failed','cancelled')
|
||||
RETURNING job_id
|
||||
"""
|
||||
),
|
||||
{"id": job_id},
|
||||
)
|
||||
.mappings()
|
||||
.first()
|
||||
)
|
||||
if row is None:
|
||||
# Ничего не обновили — либо задачи нет, либо статус неподходящий. Читаем
|
||||
# текущий статус ДО commit'а, чтобы ответ объяснял отказ, а не молчал.
|
||||
current = db.execute(
|
||||
text("SELECT status FROM nspd_geo_jobs WHERE job_id = :id"),
|
||||
{"id": job_id},
|
||||
).scalar()
|
||||
db.commit()
|
||||
return {
|
||||
"job_id": job_id,
|
||||
"resumed": False,
|
||||
"status": current,
|
||||
"reason": (
|
||||
"задача не найдена"
|
||||
if current is None
|
||||
else f"статус {current!r} не подлежит возобновлению"
|
||||
),
|
||||
}
|
||||
|
||||
db.commit()
|
||||
geo_queue = get_setting_value("nspd_geo", "queue_name", "geo")
|
||||
process_nspd_geo_job.apply_async(args=[job_id], queue=geo_queue)
|
||||
return {"job_id": job_id, "resumed": True}
|
||||
return {"job_id": job_id, "resumed": True, "status": "queued"}
|
||||
|
||||
|
||||
# ── Newbuilding cross-load ETL (#976) ────────────────────────────────────────
|
||||
|
|
|
|||
107
backend/tests/api/v1/test_2464_resume_geo_job_guard.py
Normal file
107
backend/tests/api/v1/test_2464_resume_geo_job_guard.py
Normal file
|
|
@ -0,0 +1,107 @@
|
|||
"""#2464: resume_geo_job обновлял статус БЕЗ фильтра — в отличие от cancel_geo_job.
|
||||
|
||||
`cancel_geo_job` строкой выше фильтрует явно:
|
||||
|
||||
WHERE job_id = :id AND status IN ('queued','running','paused')
|
||||
|
||||
а resume не фильтровал вовсе:
|
||||
|
||||
UPDATE nspd_geo_jobs SET status='queued', error=NULL WHERE job_id=:id
|
||||
|
||||
Следствия: завершённую задачу (done) можно было перевести обратно в queued и
|
||||
прогнать заново, затирая результат; уже бегущую — поставить в очередь второй раз,
|
||||
получив двух воркеров на один job_id и лишние запросы к НСПД, у которого WAF.
|
||||
|
||||
Замер на проде 19.08: все 66 задач в терминальных статусах (61 done, 5 cancelled).
|
||||
То есть resume на ЛЮБУЮ существующую делал ровно то, чего не должен.
|
||||
|
||||
Второе: ручка возвращала resumed=True всегда, независимо от того, изменилось ли
|
||||
что-нибудь.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
from app.api.v1 import admin_scrape
|
||||
|
||||
|
||||
class _Db:
|
||||
"""Сессия-двойник: помнит SQL и отдаёт статус задачи."""
|
||||
|
||||
def __init__(self, status: str | None, *, update_matches: bool) -> None:
|
||||
self.status = status
|
||||
self._update_matches = update_matches
|
||||
self.sql_seen: list[str] = []
|
||||
|
||||
def execute(self, statement: Any, params: Any = None) -> MagicMock:
|
||||
sql = " ".join(str(statement).split())
|
||||
self.sql_seen.append(sql)
|
||||
r = MagicMock()
|
||||
if sql.startswith("UPDATE"):
|
||||
r.mappings.return_value.first.return_value = (
|
||||
{"job_id": 1} if self._update_matches else None
|
||||
)
|
||||
else:
|
||||
r.scalar.return_value = self.status
|
||||
return r
|
||||
|
||||
def commit(self) -> None:
|
||||
pass
|
||||
|
||||
|
||||
def _resume(db: _Db) -> dict[str, Any]:
|
||||
with patch("app.workers.tasks.nspd_geo.process_nspd_geo_job") as task:
|
||||
db.task = task # type: ignore[attr-defined]
|
||||
return admin_scrape.resume_geo_job(job_id=1, db=db) # type: ignore[arg-type]
|
||||
|
||||
|
||||
def test_update_is_guarded_by_status() -> None:
|
||||
"""UPDATE обязан нести фильтр статуса — как у соседнего cancel_geo_job."""
|
||||
db = _Db("paused", update_matches=True)
|
||||
_resume(db)
|
||||
|
||||
upd = next(s for s in db.sql_seen if s.startswith("UPDATE"))
|
||||
assert "status IN" in upd, "UPDATE без фильтра статуса — возобновляется что угодно"
|
||||
|
||||
|
||||
def test_finished_job_is_not_resumed_and_says_so() -> None:
|
||||
"""done → не возобновляем и отвечаем честно, а не resumed=True."""
|
||||
db = _Db("done", update_matches=False)
|
||||
|
||||
out = _resume(db)
|
||||
|
||||
assert out["resumed"] is False
|
||||
assert out["status"] == "done"
|
||||
assert "не подлежит возобновлению" in out["reason"]
|
||||
|
||||
|
||||
def test_missing_job_reports_absence(monkeypatch) -> None:
|
||||
"""Задачи нет → resumed=False с внятной причиной, а не тихий True."""
|
||||
db = _Db(None, update_matches=False)
|
||||
|
||||
out = _resume(db)
|
||||
|
||||
assert out["resumed"] is False
|
||||
assert out["status"] is None
|
||||
assert "не найдена" in out["reason"]
|
||||
|
||||
|
||||
def test_paused_job_is_resumed() -> None:
|
||||
"""Контроль: законный случай по-прежнему работает.
|
||||
|
||||
Проверяется ТОЛЬКО resumed — новый ключ `status` здесь не трогаем, иначе тест
|
||||
падал бы и на origin/main с KeyError, то есть по причине «в ответе нет поля», а
|
||||
не «законный путь сломан». Контроль обязан быть зелёным по обе стороны.
|
||||
"""
|
||||
db = _Db("paused", update_matches=True)
|
||||
|
||||
assert _resume(db)["resumed"] is True
|
||||
|
||||
|
||||
def test_response_reports_the_new_status() -> None:
|
||||
"""Ответ несёт статус, в который перешла задача (новый ключ контракта)."""
|
||||
db = _Db("paused", update_matches=True)
|
||||
|
||||
assert _resume(db)["status"] == "queued"
|
||||
Loading…
Add table
Reference in a new issue