feat(tradein/scheduler): boot-reap — прогоны предыдущего контейнера снимаются на старте (#3122) #3123

Merged
bot-backend merged 1 commit from feat/3122-boot-reap into main 2026-08-27 09:23:44 +00:00
2 changed files with 147 additions and 1 deletions

View file

@ -0,0 +1,91 @@
"""#3122: boot-reap — прогоны предыдущего контейнера снимаются на старте, а не через 6ч.
Прод-факт (27.08, после ночного офлайна #3119): 4 прогона 'running' со стартами
03:3007:08 при старте контейнера 07:57 процессы доказуемо мертвы, но
has_running_run блокировал их источники до 6-часового порогового reap'а: до
пяти часов слепоты на источник ровно после простоя.
Критерий boot-reap started_at < старт процесса (минус минута на дрейф),
пульс не участвует: ложные срабатывания класса #2702 (редкий пульс у длинных
прогонов) невозможны по построению. Маркер counters.boot_reaped=true открывает
таким зомби подхват чекпоинта (_resume_decision): у порогового zombie процесс
может быть жив (движущаяся точка), у boot-зомби гарантированно мёртв.
Красные на main по значению: reap-SQL не существует (capability), а
resume-вердикт для boot-зомби status_zombie вместо ok (по значению).
"""
from __future__ import annotations
import os
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
from datetime import UTC, datetime
from types import SimpleNamespace
from typing import Any
from unittest.mock import MagicMock
from scraper_kit.orchestration import scheduler as sched
# ── 1. Сам boot-reap: SQL-критерий по границе старта ─────────────────────────
class _ReapDb:
def __init__(self, rows: list[Any]) -> None:
self.rows = rows
self.params: dict[str, Any] | None = None
self.sql: str = ""
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
self.sql = str(stmt)
self.params = dict(params or {})
return MagicMock(fetchall=lambda: self.rows)
def commit(self) -> None:
pass
def test_boot_reap_criterion_is_started_at_not_heartbeat() -> None:
"""Критерий — started_at против границы старта процесса с минутой запаса;
heartbeat в SQL не участвует (иначе вернулись бы ложные срабатывания #2702)."""
db = _ReapDb([SimpleNamespace(id=5018, source="avito_newbuilding_sweep")])
boot = datetime(2026, 8, 27, 7, 57, tzinfo=UTC)
n = sched.reap_boot_zombies(db, boot)
assert n == 1
assert "started_at <" in db.sql and "60 seconds" in db.sql
assert "heartbeat" not in db.sql
assert "boot_reaped" in db.sql # маркер безопасного подхвата
assert db.params and db.params["boot"] == boot
# ── 2. Resume: boot-зомби подхватывается, пороговый — нет ────────────────────
def _candidate(status: str, counters: dict[str, Any]) -> SimpleNamespace:
return SimpleNamespace(
prev_id=5040,
prev_status=status,
prev_counters=counters,
same_params=True,
age_h=2.0,
interval_days="1",
)
def test_boot_reaped_zombie_is_resumable() -> None:
"""Зомби С маркером boot_reaped — процесс гарантированно мёртв, чекпоинт
безопасен подхват. На main: status_zombie (красный по значению)."""
rid, verdict = sched._resume_decision(
_candidate("zombie", {"boot_reaped": True, "done_buckets": ["a", "b"]})
)
assert verdict["resume_reason"] == "ok", verdict
assert rid == 5040
def test_threshold_zombie_stays_rejected() -> None:
"""Пороговый зомби БЕЗ маркера — процесс может быть жив (движущаяся точка,
см. _RESUME_STATUSES) по-прежнему отказ. Инвариант обеих эр."""
rid, verdict = sched._resume_decision(_candidate("zombie", {"done_buckets": ["a", "b"]}))
assert rid is None
assert verdict["resume_reason"] == "status_zombie"

View file

@ -410,6 +410,50 @@ def reap_zombies(db: Session) -> int:
return len(rows)
def reap_boot_zombies(db: Session, boot_time: datetime) -> int:
"""Однократный boot-reap (#3122): прогон, стартовавший раньше старта СОБСТВЕННОГО
процесса планировщика, мёртв по построению его сборщик жил в предыдущем
контейнере и умер вместе с ним. Пульс в критерии не участвует вовсе, поэтому
ложные срабатывания класса #2702 (редкий пульс у 5-часовых прогонов)
невозможны: живой прогон этого процесса не может быть старше самого процесса.
Прод-цена без этого (27.08, после ночного офлайна #3119): 4 зомби блокировали
свои источники через has_running_run до 6-часового порогового reap'а — до пяти
часов слепоты на источник ровно после простоя, когда догон нужнее всего.
Минута запаса на часовой дрейф между хостом БД и приложением. Маркер
`boot_reaped: true` в counters отличает этот исход от порогового 'zombie':
у порогового процесс МОЖЕТ быть жив (reap не убивает его см. комментарий
у _RESUME_STATUSES), у boot-зомби гарантированно мёртв, поэтому его
чекпоинт безопасен для подхвата (_resume_decision пускает 'zombie' только
с этим маркером).
"""
result = db.execute(
text(
"""
UPDATE scrape_runs
SET status = 'zombie',
finished_at = clock_timestamp(),
counters = COALESCE(counters, CAST('{}' AS jsonb))
|| CAST('{"boot_reaped": true}' AS jsonb)
WHERE status = 'running'
AND started_at < CAST(:boot AS timestamptz) - CAST('60 seconds' AS interval)
RETURNING id, source
"""
),
{"boot": boot_time},
)
rows = result.fetchall()
db.commit()
if rows:
logger.warning(
"scheduler: boot-reap — %d прогонов предыдущего контейнера сняты с 'running': %s",
len(rows),
[(r.id, r.source) for r in rows],
)
return len(rows)
def _claim_run(db: Session, schedule_row: dict[str, Any], ctx: SchedulerContext) -> int | None:
"""Claim run: INSERT scrape_runs + UPDATE scrape_schedules next_run_at.
@ -518,6 +562,10 @@ def _claim_run(db: Session, schedule_row: dict[str, Any], ctx: SchedulerContext)
# counters), а причина отказа («наш баг») ничего не говорит о полноте УЖЕ записанных
# корзин — они записаны тем же heartbeat'ом, что и у banned.
_RESUME_STATUSES = frozenset({"banned", "cancelled", "failed"})
# 'zombie' допускается к подхвату ТОЛЬКО с маркером counters.boot_reaped (#3122):
# у boot-зомби процесс гарантированно мёртв (жил в предыдущем контейнере), значит
# «движущейся точки» — причины, по которой 'zombie' исключён выше, — быть не может.
# Пороговый zombie без маркера по-прежнему отвергается.
# Длина цепочки возобновлений. Не круглое число: STALE_DIGEST_INTERVAL_FACTOR (=3) —
# уже существующий в этом файле порог «источник не собирал дольше 3× своего такта =
@ -583,7 +631,8 @@ def _resume_decision(row: Any) -> tuple[int | None, dict[str, Any]]:
"resume_chain": 0,
}
if row.prev_status not in _RESUME_STATUSES:
_boot_reaped_zombie = row.prev_status == "zombie" and prev_counters.get("boot_reaped") is True
if row.prev_status not in _RESUME_STATUSES and not _boot_reaped_zombie:
verdict["resume_reason"] = f"status_{row.prev_status}"
elif not row.same_params:
verdict["resume_reason"] = "params_changed"
@ -1035,8 +1084,10 @@ async def scheduler_loop(ctx: SchedulerContext, registry: Mapping[str, Handler])
scheduler'у; 27-веточный if/elif заменён единственным `resolve_handler` + `_dispatch`.
"""
logger.info("scheduler: started (tick=%ds)", SCHEDULER_TICK_SEC)
boot_time = datetime.now(UTC) # граница «моих» прогонов для boot-reap (#3122)
# Initial sleep 30s чтобы дать FastAPI startup завершиться
await asyncio.sleep(30)
_boot_reap_done = False
while True:
# #1182 Phase 2: кооперативный SIGTERM-drain. Не reap'аем и не claim'аем
# новые run'ы во время shutdown — даём текущему dispatch'у докатиться и выходим.
@ -1046,6 +1097,10 @@ async def scheduler_loop(ctx: SchedulerContext, registry: Mapping[str, Handler])
try:
db = ctx.session_factory()
try:
# #3122: однократно на старте — снять прогоны предыдущего контейнера.
if not _boot_reap_done:
reap_boot_zombies(db, boot_time)
_boot_reap_done = True
# Reap zombies first
reap_zombies(db)
# #2670: раз в сутки — сводка «что сейчас не собирает». Календарная, а