fix(tradein/scrape_runs): оборванный деплоем прогон вне лестниц стриков и honest-status (#3393)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m57s
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m57s
После #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 — по ней и идёт отбор.
This commit is contained in:
parent
aeb1b2b3b9
commit
5f94bf0eeb
3 changed files with 274 additions and 30 deletions
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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"]
|
||||
|
|
@ -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(
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue