fix(ptica): на worker_ready зомби помечается ЛЮБОЙ 'running', а не только со снапшотом (#2464) #2975
2 changed files with 208 additions and 18 deletions
|
|
@ -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")
|
||||
|
|
|
|||
157
backend/tests/workers/test_2464_zombie_any_running.py
Normal file
157
backend/tests/workers/test_2464_zombie_any_running.py
Normal 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}"
|
||||
Loading…
Add table
Reference in a new issue