From 5f94bf0eeb0789d36632abab52f3fb4fcbebbbf2 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sun, 6 Sep 2026 09:42:27 +0500 Subject: [PATCH] =?UTF-8?q?fix(tradein/scrape=5Fruns):=20=D0=BE=D0=B1?= =?UTF-8?q?=D0=BE=D1=80=D0=B2=D0=B0=D0=BD=D0=BD=D1=8B=D0=B9=20=D0=B4=D0=B5?= =?UTF-8?q?=D0=BF=D0=BB=D0=BE=D0=B5=D0=BC=20=D0=BF=D1=80=D0=BE=D0=B3=D0=BE?= =?UTF-8?q?=D0=BD=20=D0=B2=D0=BD=D0=B5=20=D0=BB=D0=B5=D1=81=D1=82=D0=BD?= =?UTF-8?q?=D0=B8=D1=86=20=D1=81=D1=82=D1=80=D0=B8=D0=BA=D0=BE=D0=B2=20?= =?UTF-8?q?=D0=B8=20honest-status=20(#3393)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit После #3392 SIGTERM-дрейн финализирует in-flight прогоны штатными mark_done/mark_failed с counters.interrupted=1 — раньше они оставались 'running' → 'zombie' и сторожей не касались. В популяции стриков эти строки судят частичные счётчики: - detail-бэкфилл, убитый на 20% отказов, получал 'failed' с диагнозом «сбор деградировал» (_failed_ratio_too_high) — диагноз про площадку, которого никто не измерял, — и входил в стрик неудач; - cian-бэкфилл (без attempted) получал 'done' и ОБНУЛЯЛ стрик банов — ровно вред, задокументированный в orchestration/scheduler.py:99 («5 банов подряд обнулил один 'cancelled' 09.08»). Прогон с меткой interrupted теперь считается так, будто его не было: он стрик ни продлевает, ни обнуляет (обе лестницы, обе копии модуля), а три honest-status-гейта в mark_done пропускаются — статус остаётся 'done' с меткой interrupted (конвенция #3319/#3333/#3355), причина в логе. 'cancelled' оставлен как был: там прогон прервал человек. SELECT сторожа неудач дополнен колонкой counters — по ней и идёт отбор. --- .../backend/app/services/scrape_runs.py | 60 ++++-- ...t_3393_interrupted_runs_outside_streaks.py | 184 ++++++++++++++++++ .../src/scraper_kit/orchestration/runs.py | 60 ++++-- 3 files changed, 274 insertions(+), 30 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_3393_interrupted_runs_outside_streaks.py diff --git a/tradein-mvp/backend/app/services/scrape_runs.py b/tradein-mvp/backend/app/services/scrape_runs.py index ba5d65d8..3df5c256 100644 --- a/tradein-mvp/backend/app/services/scrape_runs.py +++ b/tradein-mvp/backend/app/services/scrape_runs.py @@ -102,6 +102,21 @@ def _leading_streak(rows: list[Any], is_bad: Callable[[Any], bool]) -> int: return streak +def _was_interrupted(counters: Mapping[str, Any] | None) -> bool: + """Прогон оборван SIGTERM-дрейном (#3391/#3363): counters ЧАСТИЧНЫ по построению. + + До #3392 деплой оставлял такую строку в 'running' → 'zombie', и сторожей она не + касалась. Теперь дрейн финализирует её штатным `mark_done`/`mark_failed` с + `interrupted=1` — и в лестницах стриков (#3393) её надо считать так, как будто + прогона не было: судить по частичным счётчикам нельзя ни в одну сторону, поэтому + он стрик не продлевает и не обнуляет. Обнуление и есть задокументированный вред + (kit `orchestration/scheduler.py:99`: «5 банов подряд обнулил один "cancelled" + 09.08»). 'cancelled' оставлен как был: там прогон прервал человек, и разрыв + стрика — осознанное решение, а не побочный эффект деплоя. + """ + return bool((counters or {}).get("interrupted")) + + # #2686: диагноз оборванного прогона. Пишется в scrape_runs.ban_kind (миграция 218) # РЯДОМ со status='banned', а не ВМЕСТО него — сознательный выбор между «новый # статус» и «явное поле причины»: @@ -435,7 +450,7 @@ def _alert_if_consecutive_failures(db: Session, source: str) -> None: rows = db.execute( text( """ - SELECT status FROM scrape_runs + SELECT status, counters FROM scrape_runs WHERE source = :source AND status IN ('failed', 'banned', 'done', 'cancelled') ORDER BY finished_at DESC NULLS LAST @@ -445,6 +460,10 @@ def _alert_if_consecutive_failures(db: Session, source: str) -> None: {"source": source, "limit": STREAK_SCAN_LIMIT}, ).fetchall() + # #3393: оборванного деплоем прогона в популяции стрика нет. getattr — + # функция целиком под `except Exception: pass`, и строка без колонки + # выключила бы сторож молча (тот же класс, что #2703). + rows = [r for r in rows if not _was_interrupted(getattr(r, "counters", None))] streak = _leading_streak(rows, lambda r: r.status in ("failed", "banned")) capped = streak >= STREAK_SCAN_LIMIT if not capped and not _streak_alert_due(streak, n): @@ -505,6 +524,9 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None: {"source": source, "limit": STREAK_SCAN_LIMIT}, ).fetchall() + # #3393: та же популяция, что у сторожа неудач — оборванный деплоем прогон + # не судим (в т.ч. не он решает, «измерен» ли свежайший результат ниже). + rows = [r for r in rows if not _was_interrupted(getattr(r, "counters", None))] if not rows: return @@ -692,21 +714,29 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None: даже если собрано > 0 (см. _failed_ratio_too_high). Отличие от #2625/#2700: те два смотрят на «всё или ничего» (все якоря / вся фаза), этот — на ДОЛЮ отказов у detail-backfill'ов, где ни один из первых двух признаков не матчит форму counters. + + #3393: все три гейта пропускаются, если прогон оборван дрейном (`interrupted=1`). + Их вход — счётчики ЗАВЕРШЁННОЙ работы; у оборванного они частичны, и диагноз про + площадку («сбор деградировал») был бы выдуман: бэкфилл, убитый деплоем на 20% + отказов, — это про деплой, а не про площадку. Статус остаётся штатным 'done' с + меткой interrupted (конвенция #3319/#3333/#3355), причина — в логе. """ - did_nothing = _sweep_run_did_nothing(counters) - if did_nothing is not None: - logger.error("%s run_id=%d", did_nothing, run_id) - mark_failed(db, run_id, did_nothing, counters) - return - phase_dead = _phase_totally_failed(counters) - if phase_dead is not None: - logger.error("%s run_id=%d", phase_dead, run_id) - mark_failed(db, run_id, phase_dead, counters) - return - ratio_bad = _failed_ratio_too_high(counters) - if ratio_bad is not None: - logger.error("%s run_id=%d", ratio_bad, run_id) - mark_failed(db, run_id, ratio_bad, counters) + if _was_interrupted(counters): + logger.warning( + "mark_done: run_id=%d оборван деплоем (SIGTERM-drain), частичный результат — " + "honest-status-гейты пропущены (#3393)", + run_id, + ) + honest_bad = None + else: + honest_bad = ( + _sweep_run_did_nothing(counters) + or _phase_totally_failed(counters) + or _failed_ratio_too_high(counters) + ) + if honest_bad is not None: + logger.error("%s run_id=%d", honest_bad, run_id) + mark_failed(db, run_id, honest_bad, counters) return total_seen, new_count = _column_counts(counters) row = db.execute( diff --git a/tradein-mvp/backend/tests/test_3393_interrupted_runs_outside_streaks.py b/tradein-mvp/backend/tests/test_3393_interrupted_runs_outside_streaks.py new file mode 100644 index 00000000..c4fd4042 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3393_interrupted_runs_outside_streaks.py @@ -0,0 +1,184 @@ +"""#3393: прогон, оборванный деплоем, не участвует в лестницах стриков и не судится +honest-status-гейтами. + +После #3392 SIGTERM-дрейн финализирует in-flight прогоны штатными mark_done/mark_failed +с `counters.interrupted=1` (раньше они оставались 'running' → 'zombie' и сторожей не +касались). Два следствия, которые здесь и закрыты: + + 1. Лестницы стриков увидели эти строки. У cian-бэкфиллов дрейн даёт 'done' — + и он ОБНУЛЯЛ стрик банов источника; ровно этот вред задокументирован в + kit `orchestration/scheduler.py:99` («5 банов подряд обнулил один 'cancelled' + 09.08»). У detail-бэкфиллов с counters вида attempted/failed дрейн даёт 'failed' + (см. п. 2) — и, наоборот, ПРОДЛЕВАЛ стрик. Семантика фикса: прогона как будто не + было — он стрик ни продлевает, ни обнуляет. + + 2. `_failed_ratio_too_high` судила ЧАСТИЧНЫЕ counters: бэкфилл, убитый деплоем на 20% + отказов, получал 'failed' с текстом «сбор деградировал…» — диагноз про площадку, + которого никто не измерял, плюс вход в сторож неудач. + +Проверяем на ОБЕИХ копиях (kit + app) — тем же паттерном, что +test_2670_streak_and_partial_coverage.py. +""" + +from __future__ import annotations + +import os +from types import SimpleNamespace +from typing import Any +from unittest.mock import MagicMock, patch + +import pytest + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +from scraper_kit.orchestration import runs as kit_runs + +from app.services import scrape_runs as app_runs + +_MODULES = {"kit": kit_runs, "app": app_runs} + + +def _db(rows: list[Any]) -> MagicMock: + db = MagicMock() + db.execute.return_value.fetchall.return_value = rows + return db + + +def _row(status: str, **counters: Any) -> SimpleNamespace: + return SimpleNamespace(status=status, counters=counters) + + +def _run_failure_watchdog(mod: Any, rows: list[SimpleNamespace]) -> MagicMock: + sentry = MagicMock() + with patch.object(mod, "sentry_sdk", sentry): + mod._alert_if_consecutive_failures(_db(rows), "cian_history_backfill") + return sentry + + +def _run_zero_watchdog(mod: Any, rows: list[SimpleNamespace]) -> MagicMock: + sentry = MagicMock() + with patch.object(mod, "sentry_sdk", sentry): + mod._alert_if_consecutive_zero_results(_db(rows), "cian_city_sweep") + return sentry + + +# ── (а) 'done' дрейна не ОБНУЛЯЕТ стрик банов ──────────────────────────────── + + +@pytest.mark.parametrize("name", list(_MODULES)) +def test_drained_done_does_not_reset_the_ban_streak(name: str) -> None: + """banned, banned, done+interrupted, banned → стрик 3 (веха), а не 1. + + Прод-форма из scheduler.py:99: серию банов рвал один прогон, убитый деплоем. + Красный на main: 'done' в середине обрывает `_leading_streak` на 2 → алерта нет. + """ + rows = [ + _row("banned"), + _row("banned"), + _row("done", lots_fetched=17, interrupted=1), + _row("banned"), + _row("done", lots_fetched=42), + ] + sentry = _run_failure_watchdog(_MODULES[name], rows) + sentry.capture_message.assert_called_once() + assert "3 consecutive" in sentry.capture_message.call_args[0][0] + + +# ── (б) 'failed' дрейна не ПРОДЛЕВАЕТ стрик ────────────────────────────────── + + +@pytest.mark.parametrize("name", list(_MODULES)) +def test_drained_failed_does_not_extend_the_streak(name: str) -> None: + """failed, failed, failed+interrupted, failed → стрик 3, а не 4. + + Симметрия к (а): «как будто прогона не было» работает в обе стороны. Разница + видима именно на вехе — 3 попадает в лестницу напоминаний, 4 в неё не попадает. + Красный на main: прерванный считается четвёртым → `_streak_alert_due(4, 3)` + False → сторож молчит ровно там, где отказов на самом деле три подряд. + """ + rows = [ + _row("failed"), + _row("failed"), + _row("failed", attempted=10, failed=2, interrupted=1), + _row("failed"), + _row("done", lots_fetched=42), + ] + sentry = _run_failure_watchdog(_MODULES[name], rows) + sentry.capture_message.assert_called_once() + assert "3 consecutive" in sentry.capture_message.call_args[0][0] + + +@pytest.mark.parametrize("name", list(_MODULES)) +def test_interrupted_run_is_invisible_to_the_zero_result_watchdog(name: str) -> None: + """Та же семантика у второй лестницы: zero-стрик 3, а не 4 (веха пройдена).""" + rows = [ + _row("done", lots_fetched=0), + _row("done", lots_fetched=0), + _row("done", lots_fetched=0, interrupted=1), + _row("done", lots_fetched=0), + _row("done", lots_fetched=42), + ] + sentry = _run_zero_watchdog(_MODULES[name], rows) + sentry.capture_message.assert_called_once() + assert "3 consecutive" in sentry.capture_message.call_args[0][0] + + +@pytest.mark.parametrize("name", list(_MODULES)) +def test_streaks_unchanged_without_the_mark(name: str) -> None: + """Контроль: без метки ничего не меняется — 4 отказа подряд по-прежнему между + вехами (молчание), т.е. тесты выше краснеют из-за метки, а не из-за формы строк.""" + rows = [_row("failed") for _ in range(4)] + [_row("done", lots_fetched=42)] + _run_failure_watchdog(_MODULES[name], rows).capture_message.assert_not_called() + + +# ── (в) honest-status-гейт не судит частичные counters ─────────────────────── + + +def _capture_status(mod: Any, counters: dict[str, Any]) -> list[str]: + """Статусы всех UPDATE'ов mark_done на фейковой сессии (helper из + test_honest_run_status_failed_ratio.py — читаем СТАТУС В SQL, а не имя функции).""" + statuses: list[str] = [] + + def _execute(stmt: Any, *args: Any, **kwargs: Any) -> MagicMock: + sql = str(stmt) + for status in ("done", "failed", "banned"): + if f"status = '{status}'" in sql: + statuses.append(status) + return MagicMock() + + db = MagicMock() + db.execute.side_effect = _execute + with patch.object(mod, "sentry_sdk", MagicMock()): + mod.mark_done(db, 1, dict(counters)) + return statuses + + +@pytest.mark.parametrize("name", list(_MODULES)) +def test_interrupted_backfill_stays_done(name: str) -> None: + """attempted=10, failed=2, interrupted=1 → 'done'. + + Красный на main: ratio 0.2 ≥ FAILED_RATIO_DEGRADED_THRESHOLD → 'failed' с текстом + «сбор деградировал», хотя мерили работу, которую оборвали на середине. + """ + counters = {"attempted": 10, "failed": 2, "enriched": 8, "interrupted": 1} + assert _capture_status(_MODULES[name], counters) == ["done"] + + +@pytest.mark.parametrize("name", list(_MODULES)) +def test_interrupted_sweep_with_dead_phase_stays_done(name: str) -> None: + """Тот же пропуск у #2625/#2700-гейтов: у оборванного прогона «все якоря отказали» + и «фаза мертва» — это тоже про деплой, а не про площадку.""" + counters = {"anchors_total": 4, "errors_count": 4, "lots_fetched": 0, "interrupted": 1} + assert _capture_status(_MODULES[name], counters) == ["done"] + + +@pytest.mark.parametrize("name", list(_MODULES)) +def test_same_counters_without_the_mark_still_fail(name: str) -> None: + """Контроль к обоим предыдущим: без `interrupted` гейты работают как раньше + (#3319/#3333/#3355 не отменяются, ослаблен ровно один случай).""" + assert _capture_status(_MODULES[name], {"attempted": 10, "failed": 2, "enriched": 8}) == [ + "failed" + ] + assert _capture_status( + _MODULES[name], {"anchors_total": 4, "errors_count": 4, "lots_fetched": 0} + ) == ["failed"] diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py index 85fcacf9..7d2c82bb 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py @@ -96,6 +96,21 @@ def _leading_streak(rows: list[Any], is_bad: Callable[[Any], bool]) -> int: return streak +def _was_interrupted(counters: Mapping[str, Any] | None) -> bool: + """Прогон оборван SIGTERM-дрейном (#3391/#3363): counters ЧАСТИЧНЫ по построению. + + До #3392 деплой оставлял такую строку в 'running' → 'zombie', и сторожей она не + касалась. Теперь дрейн финализирует её штатным `mark_done`/`mark_failed` с + `interrupted=1` — и в лестницах стриков (#3393) её надо считать так, как будто + прогона не было: судить по частичным счётчикам нельзя ни в одну сторону, поэтому + он стрик не продлевает и не обнуляет. Обнуление и есть задокументированный вред + (`orchestration/scheduler.py:99`: «5 банов подряд обнулил один "cancelled" 09.08»). + 'cancelled' оставлен как был: там прогон прервал человек, и разрыв стрика — + осознанное решение, а не побочный эффект деплоя. + """ + return bool((counters or {}).get("interrupted")) + + # #2686: диагноз оборванного прогона. Пишется в scrape_runs.ban_kind (миграция 218) # РЯДОМ со status='banned', а не ВМЕСТО него — сознательный выбор между «новый # статус» и «явное поле причины»: @@ -434,7 +449,7 @@ def _alert_if_consecutive_failures(db: Session, source: str) -> None: rows = db.execute( text( """ - SELECT status FROM scrape_runs + SELECT status, counters FROM scrape_runs WHERE source = :source AND status IN ('failed', 'banned', 'done', 'cancelled') ORDER BY finished_at DESC NULLS LAST @@ -444,6 +459,10 @@ def _alert_if_consecutive_failures(db: Session, source: str) -> None: {"source": source, "limit": STREAK_SCAN_LIMIT}, ).fetchall() + # #3393: оборванного деплоем прогона в популяции стрика нет. getattr — + # функция целиком под `except Exception: pass`, и строка без колонки + # выключила бы сторож молча (тот же класс, что #2703). + rows = [r for r in rows if not _was_interrupted(getattr(r, "counters", None))] streak = _leading_streak(rows, lambda r: r.status in ("failed", "banned")) capped = streak >= STREAK_SCAN_LIMIT if not capped and not _streak_alert_due(streak, n): @@ -506,6 +525,9 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None: {"source": source, "limit": STREAK_SCAN_LIMIT}, ).fetchall() + # #3393: та же популяция, что у сторожа неудач — оборванный деплоем прогон + # не судим (в т.ч. не он решает, «измерен» ли свежайший результат ниже). + rows = [r for r in rows if not _was_interrupted(getattr(r, "counters", None))] if not rows: return @@ -766,21 +788,29 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None: даже если собрано > 0 (см. _failed_ratio_too_high). Отличие от #2625/#2700: те два смотрят на «всё или ничего» (все якоря / вся фаза), этот — на ДОЛЮ отказов у detail-backfill'ов, где ни один из первых двух признаков не матчит форму counters. + + #3393: все три гейта пропускаются, если прогон оборван дрейном (`interrupted=1`). + Их вход — счётчики ЗАВЕРШЁННОЙ работы; у оборванного они частичны, и диагноз про + площадку («сбор деградировал») был бы выдуман: бэкфилл, убитый деплоем на 20% + отказов, — это про деплой, а не про площадку. Статус остаётся штатным 'done' с + меткой interrupted (конвенция #3319/#3333/#3355), причина — в логе. """ - did_nothing = _sweep_run_did_nothing(counters) - if did_nothing is not None: - logger.error("%s run_id=%d", did_nothing, run_id) - mark_failed(db, run_id, did_nothing, counters) - return - phase_dead = _phase_totally_failed(counters) - if phase_dead is not None: - logger.error("%s run_id=%d", phase_dead, run_id) - mark_failed(db, run_id, phase_dead, counters) - return - ratio_bad = _failed_ratio_too_high(counters) - if ratio_bad is not None: - logger.error("%s run_id=%d", ratio_bad, run_id) - mark_failed(db, run_id, ratio_bad, counters) + if _was_interrupted(counters): + logger.warning( + "mark_done: run_id=%d оборван деплоем (SIGTERM-drain), частичный результат — " + "honest-status-гейты пропущены (#3393)", + run_id, + ) + honest_bad = None + else: + honest_bad = ( + _sweep_run_did_nothing(counters) + or _phase_totally_failed(counters) + or _failed_ratio_too_high(counters) + ) + if honest_bad is not None: + logger.error("%s run_id=%d", honest_bad, run_id) + mark_failed(db, run_id, honest_bad, counters) return total_seen, new_count = _column_counts(counters) row = db.execute(