Compare commits
No commits in common. "8ad1e9f5ef595cbc0e2b43d33136ab37344591be" and "610ed2039555783a933471a5f94e6da10b32e8ad" have entirely different histories.
8ad1e9f5ef
...
610ed20395
2 changed files with 0 additions and 188 deletions
|
|
@ -242,59 +242,6 @@ def _resume_zombie_runs(sender=None, **_kwargs) -> None:
|
||||||
logger.warning("worker_ready: failed to enqueue geo resume job=%s: %s", jid, e)
|
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))
|
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).
|
# Sanity check: nspd_quarter_dumps table must exist (migration 88).
|
||||||
# Logs critical error but does NOT crash the worker — table may be absent
|
# Logs critical error but does NOT crash the worker — table may be absent
|
||||||
# in dev/staging before migration is applied.
|
# in dev/staging before migration is applied.
|
||||||
|
|
|
||||||
|
|
@ -1,135 +0,0 @@
|
||||||
"""На 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