Merge pull request 'fix(tradein): прогоны, оборванные деплоем (interrupted=1), невидимы лестницам стриков и honest-status-гейтам (#3393)' (#3395) from fix/3393-interrupted-runs-outside-streaks into main
Some checks failed
Deploy Trade-In / changes (push) Has been cancelled
Deploy Trade-In / build-backend (push) Has been cancelled
Deploy Trade-In / build-frontend (push) Has been cancelled
Deploy Trade-In / build-browser (push) Has been cancelled
Deploy Trade-In / deploy (push) Has been cancelled
Deploy Trade-In / deploy-status (push) Has been cancelled
Deploy Trade-In / test (push) Has been cancelled
Deploy Trade-In / perimeter-smoke (push) Has been cancelled

This commit is contained in:
bot-backend 2026-09-06 05:53:24 +00:00
commit fd4bd65126
3 changed files with 360 additions and 34 deletions

View file

@ -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,8 +460,20 @@ def _alert_if_consecutive_failures(db: Session, source: str) -> None:
{"source": source, "limit": STREAK_SCAN_LIMIT},
).fetchall()
# #3393: оборванного деплоем прогона в популяции стрика нет. getattr —
# функция целиком под `except Exception: pass`, и строка без колонки
# выключила бы сторож молча (тот же класс, что #2703).
scanned = len(rows)
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
# Потолок — по СКАНИРОВАННОМУ окну, а не по отфильтрованному: по второму
# одна interrupted-строка внутри 500 опускает измеримый максимум до 499, и
# `capped` становится недостижим по построению. Источник, сломанный наглухо,
# замолкал бы после 384-й неудачи навсегда (следующая веха 768 больше окна),
# а новую interrupted-строку в окно подкладывает каждый деплой — ровно
# анти-спам «один раз навсегда» из #2670/#2703. Условие читается так: окно
# было полным И стрик покрывает всё, что мы вообще могли судить.
capped = scanned >= STREAK_SCAN_LIMIT and streak >= len(rows)
if not capped and not _streak_alert_due(streak, n):
return
@ -505,6 +532,10 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None:
{"source": source, "limit": STREAK_SCAN_LIMIT},
).fetchall()
# #3393: та же популяция, что у сторожа неудач — оборванный деплоем прогон
# не судим (в т.ч. не он решает, «измерен» ли свежайший результат ниже).
scanned = len(rows)
rows = [r for r in rows if not _was_interrupted(getattr(r, "counters", None))]
if not rows:
return
@ -524,7 +555,10 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None:
return
streak = _leading_streak(rows, _is_zero_done)
capped = streak >= STREAK_SCAN_LIMIT
# Потолок — по сканированному окну, как у _alert_if_consecutive_failures:
# по отфильтрованному списку одна interrupted-строка делает `capped`
# недостижимым, и лестница за последней вехой в пределах окна молчит навсегда.
capped = scanned >= STREAK_SCAN_LIMIT and streak >= len(rows)
if not capped and not _streak_alert_due(streak, n):
return
@ -692,21 +726,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(

View file

@ -0,0 +1,242 @@
"""#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()
@pytest.mark.parametrize("name", list(_MODULES))
def test_zero_streak_unchanged_without_the_mark(name: str) -> None:
"""Тот же контроль для zero-result-сторожа: 4 измеренных нуля + 'done' с
результатом, без метки стрик 4, между вехами, молчание (прежнее поведение)."""
rows = [_row("done", lots_fetched=0) for _ in range(4)] + [_row("done", lots_fetched=42)]
_run_zero_watchdog(_MODULES[name], rows).capture_message.assert_not_called()
# ── (б2) потолок окна: одна interrupted-строка не глушит лестницу навсегда ────
def _full_window(mod: Any, row: SimpleNamespace, interrupted: SimpleNamespace | None) -> list[Any]:
"""Полное окно сканирования из одинаковых строк; `interrupted` (если дан) — внутри."""
rows = [row for _ in range(mod.STREAK_SCAN_LIMIT)]
if interrupted is not None:
rows[100] = interrupted
return rows
@pytest.mark.parametrize("name", list(_MODULES))
def test_full_failure_window_alerts_despite_one_interrupted(name: str) -> None:
"""500 отказов, одна строка с меткой внутри окна → алерт по потолку, не тишина.
Красный на main: `capped` считался по УЖЕ отфильтрованному списку, где максимум
499 то есть недостижим по построению, как только в окно попадает хоть одна
interrupted-строка (а её подкладывает каждый деплой). Стрик 499 мимо вехи
(последняя в пределах окна 384, следующая 768), поэтому источник, сломанный
наглухо, замолкал бы навсегда: ровно анти-спам «один раз навсегда» (#2670/#2703).
"""
mod = _MODULES[name]
limit = mod.STREAK_SCAN_LIMIT
marked = _row("failed", attempted=10, failed=2, interrupted=1)
sentry = _run_failure_watchdog(mod, _full_window(mod, _row("failed"), marked))
sentry.capture_message.assert_called_once()
assert f"{limit - 1} consecutive" in sentry.capture_message.call_args[0][0]
# Контроль: то же окно без метки алертило и до правки — краснеет метка, а не длина.
control = _run_failure_watchdog(mod, _full_window(mod, _row("failed"), None))
control.capture_message.assert_called_once()
assert f"{limit} consecutive" in control.capture_message.call_args[0][0]
@pytest.mark.parametrize("name", list(_MODULES))
def test_full_zero_window_alerts_despite_one_interrupted(name: str) -> None:
"""Тот же потолок у второй лестницы — обе копии правятся одинаково."""
mod = _MODULES[name]
limit = mod.STREAK_SCAN_LIMIT
zero = _row("done", lots_fetched=0)
marked = _row("done", lots_fetched=0, interrupted=1)
sentry = _run_zero_watchdog(mod, _full_window(mod, zero, marked))
sentry.capture_message.assert_called_once()
assert f"{limit - 1} consecutive" in sentry.capture_message.call_args[0][0]
control = _run_zero_watchdog(mod, _full_window(mod, zero, None))
control.capture_message.assert_called_once()
assert f"{limit} consecutive" in control.capture_message.call_args[0][0]
# ── (в) 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"]

View file

@ -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,8 +459,20 @@ def _alert_if_consecutive_failures(db: Session, source: str) -> None:
{"source": source, "limit": STREAK_SCAN_LIMIT},
).fetchall()
# #3393: оборванного деплоем прогона в популяции стрика нет. getattr —
# функция целиком под `except Exception: pass`, и строка без колонки
# выключила бы сторож молча (тот же класс, что #2703).
scanned = len(rows)
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
# Потолок — по СКАНИРОВАННОМУ окну, а не по отфильтрованному: по второму
# одна interrupted-строка внутри 500 опускает измеримый максимум до 499, и
# `capped` становится недостижим по построению. Источник, сломанный наглухо,
# замолкал бы после 384-й неудачи навсегда (следующая веха 768 больше окна),
# а новую interrupted-строку в окно подкладывает каждый деплой — ровно
# анти-спам «один раз навсегда» из #2670/#2703. Условие читается так: окно
# было полным И стрик покрывает всё, что мы вообще могли судить.
capped = scanned >= STREAK_SCAN_LIMIT and streak >= len(rows)
if not capped and not _streak_alert_due(streak, n):
return
@ -506,6 +533,10 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None:
{"source": source, "limit": STREAK_SCAN_LIMIT},
).fetchall()
# #3393: та же популяция, что у сторожа неудач — оборванный деплоем прогон
# не судим (в т.ч. не он решает, «измерен» ли свежайший результат ниже).
scanned = len(rows)
rows = [r for r in rows if not _was_interrupted(getattr(r, "counters", None))]
if not rows:
return
@ -525,7 +556,10 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None:
return
streak = _leading_streak(rows, _is_zero_done)
capped = streak >= STREAK_SCAN_LIMIT
# Потолок — по сканированному окну, как у _alert_if_consecutive_failures:
# по отфильтрованному списку одна interrupted-строка делает `capped`
# недостижимым, и лестница за последней вехой в пределах окна молчит навсегда.
capped = scanned >= STREAK_SCAN_LIMIT and streak >= len(rows)
if not capped and not _streak_alert_due(streak, n):
return
@ -766,21 +800,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(