fix(ptica): осиротевшие прогоны Объектива закрываются, а не висят вечно (#2464) #2978
2 changed files with 188 additions and 0 deletions
|
|
@ -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.
|
||||
|
|
|
|||
135
backend/tests/workers/test_2464_objective_zombie_sweep.py
Normal file
135
backend/tests/workers/test_2464_objective_zombie_sweep.py
Normal file
|
|
@ -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}"
|
||||
Loading…
Add table
Reference in a new issue