fix(ptica): resume_geo_job больше не возобновляет что попало (#2464) (#2946)
Some checks failed
Deploy / build-backend (push) Blocked by required conditions
Deploy / build-worker (push) Blocked by required conditions
Deploy / build-frontend (push) Blocked by required conditions
Deploy / deploy (push) Blocked by required conditions
Deploy / deploy-caddy (push) Blocked by required conditions
Deploy / perimeter-smoke (push) Blocked by required conditions
Deploy / deploy-status (push) Blocked by required conditions
Deploy / changes (push) Has been cancelled

This commit is contained in:
bot-backend 2026-08-20 06:59:55 +00:00
parent 251eacc4b6
commit 9e4b190303
3 changed files with 239 additions and 16 deletions

View file

@ -1123,18 +1123,48 @@ def cancel_geo_job(
db: Annotated[Session, Depends(get_db)],
) -> dict[str, Any]:
"""Пометить job как cancelled. Worker увидит при следующей итерации."""
db.execute(
text(
"""
UPDATE nspd_geo_jobs SET status = 'cancelled', finished_at = NOW(),
error = COALESCE(error, 'cancelled by admin')
WHERE job_id = :id AND status IN ('queued','running','paused')
"""
),
{"id": job_id},
# #2464: фильтр статуса здесь был всегда (в отличие от resume ниже), но ответ
# возвращал cancelled=True независимо от того, задел ли UPDATE хоть одну строку.
# Несуществующий job_id и уже завершённая задача давали тот же ответ, что
# настоящая отмена — оператор и админ-UI получали подтверждение действия,
# которого не было.
#
# Обоснование держим в КОММЕНТАРИИ, а не в докстринге: FastAPI кладёт докстринг
# в OpenAPI-description, откуда он попадает в опубликованный контракт и в
# сгенерированные типы фронта (frontend/src/lib/api-types.ts). Внутренние замеры
# там не нужны, а gate openapi-codegen-check честно ловит такое расхождение.
row = (
db.execute(
text(
"""
UPDATE nspd_geo_jobs SET status = 'cancelled', finished_at = NOW(),
error = COALESCE(error, 'cancelled by admin')
WHERE job_id = :id AND status IN ('queued','running','paused')
RETURNING job_id
"""
),
{"id": job_id},
)
.mappings()
.first()
)
if row is None:
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,
"cancelled": False,
"status": current,
"reason": (
"задача не найдена" if current is None else f"статус {current!r} уже терминальный"
),
}
db.commit()
return {"job_id": job_id, "cancelled": True}
return {"job_id": job_id, "cancelled": True, "status": "cancelled"}
@router.post("/geo/jobs/{job_id}/resume")
@ -1142,18 +1172,61 @@ 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,
# который фильтрует явно. Из-за этого «возобновить» можно было завершённую задачу
# (done → снова queued и повторный прогон, затирая результат) и уже бегущую
# (второй worker на тот же job_id — лишние запросы к НСПД, у которого WAF).
#
# Замер на проде 19.08: все 66 задач в терминальных статусах — 61 done, 5
# cancelled. То есть resume на ЛЮБУЮ существующую делал ровно то, чего не должен.
#
# Второе: ручка возвращала resumed=True всегда, независимо от того, изменилось ли
# что-нибудь. Теперь ответ отражает факт — статус и причина в ответе, задача НЕ
# ставится в очередь.
#
# 'cancelled' оставлен возобновляемым намеренно: cancel — ручное действие
# оператора, и без этого отменённая по ошибке задача не восстанавливалась бы.
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) ────────────────────────────────────────

View file

@ -0,0 +1,150 @@
"""#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"
# ── cancel_geo_job: тот же класс — подтверждение действия, которого не было ─────
#
# Фильтр статуса здесь был всегда, но ответ возвращал cancelled=True независимо от
# того, задел ли UPDATE строку. Несуществующий job_id и уже завершённая задача давали
# тот же ответ, что настоящая отмена.
def _cancel(db: _Db) -> dict[str, Any]:
return admin_scrape.cancel_geo_job(job_id=1, db=db) # type: ignore[arg-type]
def test_cancel_of_finished_job_is_not_reported_as_success() -> None:
db = _Db("done", update_matches=False)
out = _cancel(db)
assert out["cancelled"] is False
assert out["status"] == "done"
assert "терминальный" in out["reason"]
def test_cancel_of_missing_job_reports_absence() -> None:
db = _Db(None, update_matches=False)
out = _cancel(db)
assert out["cancelled"] is False
assert out["status"] is None
assert "не найдена" in out["reason"]
def test_cancel_of_running_job_still_works() -> None:
"""Контроль: законная отмена по-прежнему подтверждается.
Проверяется только `cancelled` новый ключ `status` намеренно не трогаем, иначе
контроль падал бы на origin/main с KeyError, то есть «в ответе нет поля», а не
«отмена сломана».
"""
db = _Db("running", update_matches=True)
assert _cancel(db)["cancelled"] is True

View file

@ -1763,7 +1763,7 @@ export interface paths {
put?: never;
/**
* Resume Geo Job
* @description Re-enqueue paused/failed job. Resume idempotent через pending targets.
* @description Re-enqueue задачу из НЕзавершённого состояния (paused / failed / cancelled).
*/
post: operations["resume_geo_job_api_v1_admin_scrape_geo_jobs__job_id__resume_post"];
delete?: never;