fix(ptica): на worker_ready зомби помечается ЛЮБОЙ 'running', а не только со снапшотом (#2464)
Some checks failed
Deploy / deploy (push) Blocked by required conditions
Deploy / deploy-status (push) Blocked by required conditions
Deploy / changes (push) Successful in 14s
Deploy / perimeter-smoke (push) Blocked by required conditions
Deploy / build-frontend (push) Has been skipped
Deploy / deploy-caddy (push) Has been skipped
Deploy / build-worker (push) Has been cancelled
Deploy / build-backend (push) Has been cancelled
Some checks failed
Deploy / deploy (push) Blocked by required conditions
Deploy / deploy-status (push) Blocked by required conditions
Deploy / changes (push) Successful in 14s
Deploy / perimeter-smoke (push) Blocked by required conditions
Deploy / build-frontend (push) Has been skipped
Deploy / deploy-caddy (push) Has been skipped
Deploy / build-worker (push) Has been cancelled
Deploy / build-backend (push) Has been cancelled
Закрывает часть эпика #2464. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
commit
18593c3019
2 changed files with 208 additions and 18 deletions
|
|
@ -89,14 +89,24 @@ def _resume_zombie_runs(sender=None, **_kwargs) -> None:
|
||||||
db = SessionLocal()
|
db = SessionLocal()
|
||||||
ids: list[int] = []
|
ids: list[int] = []
|
||||||
try:
|
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 = (
|
rows = (
|
||||||
db.execute(
|
db.execute(
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
SELECT run_id
|
SELECT run_id, (objects_snapshot IS NOT NULL) AS resumable
|
||||||
FROM kn_scrape_runs
|
FROM kn_scrape_runs
|
||||||
WHERE status = 'running'
|
WHERE status = 'running'
|
||||||
AND objects_snapshot IS NOT NULL
|
|
||||||
ORDER BY started_at ASC
|
ORDER BY started_at ASC
|
||||||
LIMIT 20
|
LIMIT 20
|
||||||
"""
|
"""
|
||||||
|
|
@ -106,22 +116,45 @@ def _resume_zombie_runs(sender=None, **_kwargs) -> None:
|
||||||
.all()
|
.all()
|
||||||
)
|
)
|
||||||
if rows:
|
if rows:
|
||||||
ids = [int(r["run_id"]) for r in rows]
|
ids = [int(r["run_id"]) for r in rows if r["resumable"]]
|
||||||
# Помечаем найденные как 'zombie' одним апдейтом — resume создаст новые
|
orphan_ids = [int(r["run_id"]) for r in rows if not r["resumable"]]
|
||||||
# run_id со ссылкой resumed_from_run_id.
|
if ids:
|
||||||
db.execute(
|
# Помечаем как 'zombie' одним апдейтом — resume создаст новые
|
||||||
text(
|
# run_id со ссылкой resumed_from_run_id.
|
||||||
"""
|
db.execute(
|
||||||
UPDATE kn_scrape_runs
|
text(
|
||||||
SET status = 'zombie',
|
"""
|
||||||
finished_at = NOW(),
|
UPDATE kn_scrape_runs
|
||||||
error = COALESCE(error,
|
SET status = 'zombie',
|
||||||
'auto-zombie at worker_ready, resume scheduled')
|
finished_at = NOW(),
|
||||||
WHERE run_id = ANY(:ids)
|
error = COALESCE(error,
|
||||||
"""
|
'auto-zombie at worker_ready, resume scheduled')
|
||||||
),
|
WHERE run_id = ANY(:ids)
|
||||||
{"ids": 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()
|
db.commit()
|
||||||
else:
|
else:
|
||||||
logger.info("worker_ready: нет stale kn runs для resume")
|
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