From af4d2a18539b0681003af3602879309c32afa4e7 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Wed, 19 Aug 2026 22:38:25 +0500 Subject: [PATCH] =?UTF-8?q?fix(ptica):=20resume=5Fgeo=5Fjob=20=D0=B1=D0=BE?= =?UTF-8?q?=D0=BB=D1=8C=D1=88=D0=B5=20=D0=BD=D0=B5=20=D0=B2=D0=BE=D0=B7?= =?UTF-8?q?=D0=BE=D0=B1=D0=BD=D0=BE=D0=B2=D0=BB=D1=8F=D0=B5=D1=82=20=D1=87?= =?UTF-8?q?=D1=82=D0=BE=20=D0=BF=D0=BE=D0=BF=D0=B0=D0=BB=D0=BE=20(#2464)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- backend/app/api/v1/admin_scrape.py | 55 ++++++++- .../api/v1/test_2464_resume_geo_job_guard.py | 107 ++++++++++++++++++ 2 files changed, 157 insertions(+), 5 deletions(-) create mode 100644 backend/tests/api/v1/test_2464_resume_geo_job_guard.py diff --git a/backend/app/api/v1/admin_scrape.py b/backend/app/api/v1/admin_scrape.py index ffa30d0c..a8e412cc 100644 --- a/backend/app/api/v1/admin_scrape.py +++ b/backend/app/api/v1/admin_scrape.py @@ -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) ──────────────────────────────────────── diff --git a/backend/tests/api/v1/test_2464_resume_geo_job_guard.py b/backend/tests/api/v1/test_2464_resume_geo_job_guard.py new file mode 100644 index 00000000..198d8a34 --- /dev/null +++ b/backend/tests/api/v1/test_2464_resume_geo_job_guard.py @@ -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"