From 63956f089150d07b308b22288332f7e76621cc6e Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 27 Aug 2026 14:17:42 +0500 Subject: [PATCH] =?UTF-8?q?feat(tradein/scheduler):=20boot-reap=20?= =?UTF-8?q?=E2=80=94=20=D0=BF=D1=80=D0=BE=D0=B3=D0=BE=D0=BD=D1=8B=20=D0=BF?= =?UTF-8?q?=D1=80=D0=B5=D0=B4=D1=8B=D0=B4=D1=83=D1=89=D0=B5=D0=B3=D0=BE=20?= =?UTF-8?q?=D0=BA=D0=BE=D0=BD=D1=82=D0=B5=D0=B9=D0=BD=D0=B5=D1=80=D0=B0=20?= =?UTF-8?q?=D1=81=D0=BD=D0=B8=D0=BC=D0=B0=D1=8E=D1=82=D1=81=D1=8F=20=D0=BD?= =?UTF-8?q?=D0=B0=20=D1=81=D1=82=D0=B0=D1=80=D1=82=D0=B5=20(#3122)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Прод-факт 27.08 (после ночного офлайна #3119): 4 прогона 'running' со стартами до старта контейнера блокировали свои источники через has_running_run до 6-часового порогового reap'а — до пяти часов слепоты на источник ровно после простоя, когда догон нужнее всего. Критерий — started_at < старт процесса планировщика (минус минута на дрейф), пульс не участвует: ложные срабатывания класса #2702 (редкий пульс длинных прогонов) невозможны по построению — живой прогон этого процесса не может быть старше самого процесса. Маркер counters.boot_reaped=true открывает boot-зомби подхват чекпоинта (_resume_decision): у порогового zombie процесс может быть жив (движущаяся точка — причина исключения 'zombie' из _RESUME_STATUSES), у boot-зомби — гарантированно мёртв. Пороговый zombie без маркера по-прежнему отвергается (закреплено тестом-инвариантом). Co-Authored-By: Claude Opus 5 --- .../backend/tests/test_3122_boot_reap.py | 91 +++++++++++++++++++ .../scraper_kit/orchestration/scheduler.py | 57 +++++++++++- 2 files changed, 147 insertions(+), 1 deletion(-) create mode 100644 tradein-mvp/backend/tests/test_3122_boot_reap.py diff --git a/tradein-mvp/backend/tests/test_3122_boot_reap.py b/tradein-mvp/backend/tests/test_3122_boot_reap.py new file mode 100644 index 00000000..1f8e9773 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3122_boot_reap.py @@ -0,0 +1,91 @@ +"""#3122: boot-reap — прогоны предыдущего контейнера снимаются на старте, а не через 6ч. + +Прод-факт (27.08, после ночного офлайна #3119): 4 прогона 'running' со стартами +03:30–07: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" diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py index 9eaed46c..029b96c7 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py @@ -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: раз в сутки — сводка «что сейчас не собирает». Календарная, а -- 2.45.3