fix(ptica): на worker_ready зомби помечается ЛЮБОЙ 'running', а не только со снапшотом (#2464) #2975

Merged
bot-backend merged 1 commit from fix/2464-zombie-any-running into main 2026-08-20 12:22:38 +00:00
2 changed files with 208 additions and 18 deletions

View file

@ -89,14 +89,24 @@ def _resume_zombie_runs(sender=None, **_kwargs) -> None:
db = SessionLocal()
ids: list[int] = []
try:
# #2464: берём ЛЮБУЮ строку в 'running', без фильтра по objects_snapshot.
# Докстринг этой функции формулирует инвариант прямо: «by definition, on
# worker_ready ANY 'running' row is a zombie because there is no active
# worker». Фильтр ему противоречил: строка без снапшота не попадала в
# выборку и оставалась 'running' НАВСЕГДА — ровно то состояние, ради
# устранения которого функция и заводилась.
#
# Снапшот всё равно нужен — но не для пометки, а для ВОЗОБНОВЛЕНИЯ:
# resume_kn_run восстанавливает обход «using objects_snapshot». Поэтому
# помечаем зомби всех, а resume ставим только тем, кого есть чем
# возобновить. Остальные получают честную причину вместо тишины.
rows = (
db.execute(
text(
"""
SELECT run_id
SELECT run_id, (objects_snapshot IS NOT NULL) AS resumable
FROM kn_scrape_runs
WHERE status = 'running'
AND objects_snapshot IS NOT NULL
ORDER BY started_at ASC
LIMIT 20
"""
@ -106,22 +116,45 @@ def _resume_zombie_runs(sender=None, **_kwargs) -> None:
.all()
)
if rows:
ids = [int(r["run_id"]) for r in rows]
# Помечаем найденные как 'zombie' одним апдейтом — resume создаст новые
# run_id со ссылкой resumed_from_run_id.
db.execute(
text(
"""
UPDATE kn_scrape_runs
SET status = 'zombie',
finished_at = NOW(),
error = COALESCE(error,
'auto-zombie at worker_ready, resume scheduled')
WHERE run_id = ANY(:ids)
"""
),
{"ids": ids},
)
ids = [int(r["run_id"]) for r in rows if r["resumable"]]
orphan_ids = [int(r["run_id"]) for r in rows if not r["resumable"]]
if ids:
# Помечаем как 'zombie' одним апдейтом — resume создаст новые
# run_id со ссылкой resumed_from_run_id.
db.execute(
text(
"""
UPDATE kn_scrape_runs
SET status = 'zombie',
finished_at = NOW(),
error = COALESCE(error,
'auto-zombie at worker_ready, resume scheduled')
WHERE run_id = ANY(:ids)
"""
),
{"ids": ids},
)
if orphan_ids:
db.execute(
text(
"""
UPDATE kn_scrape_runs
SET status = 'zombie',
finished_at = NOW(),
error = COALESCE(error,
'auto-zombie at worker_ready, resume невозможен: '
'нет objects_snapshot')
WHERE run_id = ANY(:orphans)
"""
),
{"orphans": orphan_ids},
)
logger.warning(
"worker_ready: %d kn-прогонов без objects_snapshot помечены zombie"
" без resume — возобновлять нечем: %s",
len(orphan_ids),
orphan_ids,
)
db.commit()
else:
logger.info("worker_ready: нет stale kn runs для resume")

View file

@ -0,0 +1,157 @@
"""На worker_ready зомби помечается ЛЮБАЯ строка в 'running' (#2464).
Докстринг `_resume_zombie_runs` формулирует инвариант прямо:
No time threshold: by definition, on worker_ready ANY 'running' row is a zombie
because there is no active worker. Previously we required heartbeat stayed in
'running' status forever and required manual cancel/resume.
А запрос добавлял `AND objects_snapshot IS NOT NULL`. Строка без снапшота в выборку не
попадала и оставалась `'running'` НАВСЕГДА ровно то состояние, ради устранения которого
функция и заводилась.
Снапшот нужен, но не для пометки, а для ВОЗОБНОВЛЕНИЯ: `resume_kn_run` восстанавливает
обход «using objects_snapshot». Поэтому зомби помечаются все, а resume ставится только тем,
кого есть чем возобновить; остальные получают честную причину.
Тесты смотрят на выполненный SQL и на поставленные resume-задачи то есть на поведение.
"""
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:
"""Отдаёт заданные строки на SELECT и запоминает весь выполненный SQL."""
def __init__(self, rows: list[dict]) -> None:
self._rows = rows
self.sql: list[tuple[str, dict]] = []
self.commits = 0
def execute(self, statement: Any = None, params: Any = None, *a: Any, **kw: Any) -> _Result:
text_ = str(statement)
self.sql.append((text_, params or {}))
if "SELECT run_id" in text_ and "kn_scrape_runs" in text_:
# Двойник ОБЯЗАН соблюдать ровно тот фильтр, о котором идёт спор.
# Без этого тест краснел бы на origin/main по ложной причине: там
# строка без снапшота отсекается запросом и до пометки не доходит,
# а наивный двойник отдавал бы её всё равно — и главный тест
# («строка не помечена») проходил бы на обеих сторонах.
# Ищем именно условие в WHERE (с ведущим AND), а не подвыражение в
# SELECT: правка выносит `(objects_snapshot IS NOT NULL) AS resumable`
# в список полей, и совпадение по голой подстроке отсекало бы строки
# и у исправленной версии тоже.
if "AND objects_snapshot IS NOT NULL" in text_:
return _Result([r for r in self._rows if r["resumable"]])
return _Result(self._rows)
return _Result([])
def commit(self) -> None:
self.commits += 1
def rollback(self) -> None:
pass
def close(self) -> None:
pass
def _run(rows: list[dict]) -> tuple[_Session, list[int]]:
from app.workers import lifecycle as mod
db = _Session(rows)
enqueued: list[int] = []
resume_task = MagicMock()
resume_task.apply_async.side_effect = lambda args=None, **_kw: enqueued.append(args[0])
# SessionLocal импортируется ВНУТРИ функции (late import), поэтому подменяем
# его в модуле-источнике, а не на lifecycle.
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(resume_kn_run=resume_task)},
),
):
mod._resume_zombie_runs()
return db, enqueued
def _updated_ids(db: _Session) -> set[int]:
out: set[int] = set()
for sql, params in db.sql:
if "UPDATE kn_scrape_runs" in sql and "zombie" in sql:
for key in ("ids", "orphans"):
out.update(params.get(key) or [])
return out
def test_running_row_without_snapshot_is_marked_zombie() -> None:
"""Строка без objects_snapshot обязана быть помечена, а не остаться 'running'.
На origin/main она вообще не попадает в выборку: фильтр отсекает её,
и никакого UPDATE по ней не выполняется.
"""
db, _ = _run([{"run_id": 7, "resumable": False}])
assert 7 in _updated_ids(db), (
"прогон без снапшота не помечен zombie — останется в 'running' навсегда; "
f"выполненный SQL: {[s for s, _ in db.sql]}"
)
def test_unresumable_row_gets_no_resume_task() -> None:
"""Контроль: возобновлять нечем — resume не ставим.
Ловит «починку», которая ставила бы resume всем подряд: resume_kn_run
восстанавливает обход ИЗ objects_snapshot и без него упадёт.
"""
_, enqueued = _run([{"run_id": 7, "resumable": False}])
assert enqueued == [], f"поставлен resume для невозобновляемого прогона: {enqueued}"
def test_resumable_row_still_gets_resume() -> None:
"""Контроль: прежнее поведение для строк со снапшотом не тронуто."""
db, enqueued = _run([{"run_id": 3, "resumable": True}])
assert 3 in _updated_ids(db)
assert enqueued == [3]
def test_mixed_batch_splits_correctly() -> None:
"""Контроль: смешанная выборка — помечены все, resume только у пригодных."""
db, enqueued = _run(
[
{"run_id": 1, "resumable": True},
{"run_id": 2, "resumable": False},
{"run_id": 3, "resumable": True},
]
)
assert _updated_ids(db) == {1, 2, 3}, f"помечены не все: {_updated_ids(db)}"
assert sorted(enqueued) == [1, 3], f"resume поставлен не тем: {enqueued}"