diff --git a/backend/app/workers/lifecycle.py b/backend/app/workers/lifecycle.py index 8cade6fe..6931b4b2 100644 --- a/backend/app/workers/lifecycle.py +++ b/backend/app/workers/lifecycle.py @@ -242,6 +242,59 @@ def _resume_zombie_runs(sender=None, **_kwargs) -> None: logger.warning("worker_ready: failed to enqueue geo resume job=%s: %s", jid, e) logger.info("worker_ready: resume scan finished (geo_jobs=%d)", len(geo_resume_jobs)) + # objective_scrape_runs: тот же инвариант, что у kn — на worker_ready активных + # воркеров нет, значит любая строка в 'running' осиротела. Подметальщика у этой + # таблицы не было вовсе, и на проде 2026-08-20 висело 6 строк со статусом + # 'running' с 17.05 (94 суток), при 71 'done' и НИ ОДНОГО 'failed' — след + # отравления сессии, из-за которого _finish_run(status='failed') не мог + # записаться (починено #2972). Причина устранена, но жёсткое убийство воркера + # (редеплой, OOM) по-прежнему оставляет 'running' навсегда: у Объектива нет + # ни своего cleanup_zombies, ни snapshot'а для resume. + # + # Resume не делаем — возобновлять нечего (снапшота обхода нет), только честно + # закрываем. finished_at ставим НЕ NOW(), а по последнему признаку жизни: + # прогон, умерший 94 дня назад, не должен читаться как «завершён только что». + # Монитору свежести это безразлично в обе стороны — он считает last_success_at + # и recent_output только по status='done', а last_attempt_at/last_status — по + # started_at (см. _FRESHNESS_SOURCES в admin_scrape.py), так что зомби-строки + # в него не попадают ни одним столбцом. + db = SessionLocal() + try: + rows = ( + db.execute( + text( + """ + UPDATE objective_scrape_runs + SET status = 'zombie', + finished_at = COALESCE(heartbeat_at, started_at), + error = COALESCE(error, + 'auto-zombie at worker_ready: воркер перезапущен ' + 'во время прогона, возобновление невозможно') + WHERE status = 'running' + RETURNING run_id + """ + ) + ) + .mappings() + .all() + ) + db.commit() + if rows: + logger.info( + "worker_ready: objective_scrape_runs — помечено зомби: %s", + [int(r["run_id"]) for r in rows], + ) + else: + logger.info("worker_ready: нет осиротевших objective-прогонов") + except Exception as e: + logger.warning("worker_ready objective zombie sweep failed: %s", e) + try: + db.rollback() + except Exception: + pass + finally: + db.close() + # Sanity check: nspd_quarter_dumps table must exist (migration 88). # Logs critical error but does NOT crash the worker — table may be absent # in dev/staging before migration is applied. diff --git a/backend/tests/workers/test_2464_objective_zombie_sweep.py b/backend/tests/workers/test_2464_objective_zombie_sweep.py new file mode 100644 index 00000000..25ce10d3 --- /dev/null +++ b/backend/tests/workers/test_2464_objective_zombie_sweep.py @@ -0,0 +1,135 @@ +"""На worker_ready осиротевшие прогоны Объектива закрываются, а не висят вечно (#2464). + +`objective_scrape_runs` не подметалась ничем: `worker_ready` знал только про +`kn_scrape_runs` и `nspd_geo_jobs`. На проде 2026-08-20 в ней висело **6 строк +'running' с 17.05** (94 суток) при 71 'done' и **ни одном 'failed'** — след +отравления сессии, из-за которого `_finish_run(status='failed')` не мог +записаться (причина починена #2972). Причина устранена, но жёсткое убийство +воркера по-прежнему оставляет 'running' навсегда. + +Тесты смотрят на выполненный SQL — то есть на поведение функции, а не на её текст. +На origin/main запроса к `objective_scrape_runs` нет вовсе, поэтому головной тест +там красный по существу («строка не закрыта»), а не из-за отсутствующего символа. +""" + +from __future__ import annotations + +import os + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +from typing import Any +from unittest.mock import MagicMock, patch + + +class _Result: + def __init__(self, rows: list[dict]) -> None: + self._rows = rows + + def mappings(self) -> _Result: + return self + + def all(self) -> list[dict]: + return self._rows + + def scalar(self) -> Any: + return None + + def first(self) -> Any: + return None + + +class _Session: + def __init__(self) -> None: + self.sql: list[str] = [] + + def execute(self, statement: Any = None, params: Any = None, *a: Any, **kw: Any) -> _Result: + text_ = str(statement) + self.sql.append(text_) + if "objective_scrape_runs" in text_ and "RETURNING" in text_: + return _Result([{"run_id": 1}, {"run_id": 2}]) + return _Result([]) + + def commit(self) -> None: + pass + + def rollback(self) -> None: + pass + + def close(self) -> None: + pass + + +def _run() -> _Session: + from app.workers import lifecycle as mod + + db = _Session() + import app.core.db as core_db + + with ( + patch.object(core_db, "SessionLocal", lambda: db), + patch.dict( + "sys.modules", + { + "app.workers.tasks.scrape_kn": MagicMock(), + "app.workers.tasks.nspd_geo": MagicMock(_ZOMBIE_PAUSED_THRESHOLD="30 minutes"), + }, + ), + ): + mod._resume_zombie_runs() + return db + + +def _objective_update(db: _Session) -> str | None: + for sql in db.sql: + if "UPDATE objective_scrape_runs" in sql and "zombie" in sql: + return sql + return None + + +def test_orphan_objective_run_is_closed() -> None: + """Головной: worker_ready обязан закрыть 'running'-прогоны Объектива. + + На origin/main такого запроса нет вовсе — строки остаются 'running' навсегда. + """ + db = _run() + assert _objective_update(db) is not None, ( + "worker_ready не трогает objective_scrape_runs — осиротевший прогон " + f"останется 'running' навсегда; выполнено запросов: {len(db.sql)}" + ) + + +def test_only_running_rows_are_touched() -> None: + """Контроль от переусердствования: закрываем только 'running'. + + Без этого условия подметание переписало бы 'done'-прогоны и стёрло историю. + """ + sql = _objective_update(db := _run()) + assert sql is not None + assert "WHERE status = 'running'" in sql, f"нет сужения по статусу:\n{sql}" + assert "'done'" not in sql, f"запрос упоминает 'done' — риск задеть завершённые:\n{sql}" + assert db is not None + + +def test_finished_at_is_last_sign_of_life_not_now() -> None: + """Контроль честности времени: прогон, умерший 94 дня назад, не «завершён сейчас». + + NOW() здесь соврал бы и оператору в списке прогонов, и любому будущему + потребителю finished_at. + """ + sql = _objective_update(_run()) + assert sql is not None + assert ( + "COALESCE(heartbeat_at, started_at)" in sql + ), f"finished_at ставится не по последнему признаку жизни:\n{sql}" + assert ( + "finished_at = NOW()" not in sql + ), f"finished_at = NOW() — время завершения соврано:\n{sql}" + + +def test_kn_sweep_still_runs() -> None: + """Контроль от регресса: добавление Объектива не сломало подметание kn.""" + db = _run() + assert any( + "kn_scrape_runs" in s for s in db.sql + ), f"подметание kn_scrape_runs пропало; выполнено: {db.sql}"