fix(tradein/scraper): длинная серия неудач перестаёт замолкать, оборванный прогон — называться успехом (#2670) #2720

Merged
bot-backend merged 2 commits from fix/2670-streak-and-coverage into main 2026-08-06 10:06:58 +00:00
6 changed files with 500 additions and 72 deletions

View file

@ -55,6 +55,52 @@ CONSECUTIVE_FAILURE_ALERT_THRESHOLD = 3
# невидим для CONSECUTIVE_FAILURE_ALERT_THRESHOLD (тот считает только failed/banned).
CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD = 3
# #2670: анти-спам «один раз на стрик» безопасен ТОЛЬКО там, где стрик прерывается
# не только в принципе, но и на практике. Оба сторожа ниже слали алерт ровно на N-й
# подряд неудаче и дальше молчали навсегда — а у постоянно сломанного источника
# «дальше» длится месяцами. Прод 2026-08-06: у avito_full_load 31 неудача подряд,
# последний успешный прогон 03.07 (34 дня без сбора), алерт был ровно один — на
# третьей; у avito_full_load_exhaustive 5 подряд. Тишина при этом неотличима от
# «всё хорошо» — ровно та ловушка, из-за которой #2574 месяц выглядела как норма.
#
# Вместо «ровно N» — разреженная лестница напоминаний: N, 2N, 4N, 8N…, а дальше не
# реже, чем раз в STREAK_ALERT_MAX_PERIOD×N прогонов. Лестница по ПРОГОНАМ, а не
# «раз в сутки», потому что источники идут разным тактом: domclick_city_sweep — раз
# в день, proxy_healthcheck — раз в полчаса; календарное разрежение для одного из
# них всегда будет либо спамом, либо молчанием.
STREAK_ALERT_MAX_PERIOD = 16
# Потолок сканирования истории источника при подсчёте стрика. Достигнутый потолок
# сам по себе повод для алерта (стрик заведомо огромен) — так «замолчать навсегда»
# невозможно по построению, а не по счастливому совпадению чисел.
STREAK_SCAN_LIMIT = 500
def _streak_alert_due(streak: int, threshold: int) -> bool:
"""Достиг ли стрик очередной вехи напоминания (#2670).
True на threshold, 2×, 4×, 8× и дальше на каждом кратном
STREAK_ALERT_MAX_PERIOD×threshold. Первый алерт приходит там же, где и раньше
на N-й подряд неудаче; меняется только то, что он не последний.
"""
if streak < threshold or streak % threshold:
return False
mult = streak // threshold
if mult % STREAK_ALERT_MAX_PERIOD == 0:
return True
return mult & (mult - 1) == 0
def _leading_streak(rows: list[Any], is_bad: Callable[[Any], bool]) -> int:
"""Длина серии подряд идущих «плохих» строк с начала списка (свежие — первыми)."""
streak = 0
for row in rows:
if not is_bad(row):
break
streak += 1
return streak
# #2686: диагноз оборванного прогона. Пишется в scrape_runs.ban_kind (миграция 218)
# РЯДОМ со status='banned', а не ВМЕСТО него — сознательный выбор между «новый
# статус» и «явное поле причины»:
@ -147,20 +193,27 @@ def _column_counts(counters: dict[str, int]) -> tuple[int | None, int | None]:
def _alert_if_consecutive_failures(db: Session, source: str) -> None:
"""Отправить Sentry alert если последние CONSECUTIVE_FAILURE_ALERT_THRESHOLD
завершённых запусков для данного source имеют статус 'failed' или 'banned'.
"""Sentry alert на серию из CONSECUTIVE_FAILURE_ALERT_THRESHOLD неудач подряд
(статусы 'failed'/'banned') у данного source.
Anti-spam: алерт срабатывает ТОЛЬКО когда стрик РОВНО равен порогу т.е. запрос
возвращает ровно N последних (failed|banned) и (N+1)-й, если существует, НЕ является
failed/banned. Это предотвращает повторный алерт на каждой ошибке сверх порога.
Anti-spam: не на каждой неудаче, а по разреженной лестнице вех (см.
_streak_alert_due). До #2670 алерт приходил РОВНО на N-й неудаче и дальше не
повторялся никогда: серия, ставшая длиннее порога, замолкала навсегда. На проде
это дало avito_full_load 31 неудача подряд, 34 дня без единого успешного
прогона, один алерт за всё время.
Стрик прерывается любым завершением, кроме failed/banned, по данным прода это
достижимо и достигается (у domclick_city_sweep текущий стрик равен 1 при 47
завершённых прогонах), поэтому лестница не вырождается в постоянный алерт.
Best-effort: весь блок обёрнут в try/except сбой запроса или неинициализированный
Sentry НЕ должен нарушать вызывающий mark_* путь.
"""
if sentry_sdk is None:
return
n = CONSECUTIVE_FAILURE_ALERT_THRESHOLD
try:
# Берём последние N+1 завершённых (non-running) запусков по source.
# Сортируем по finished_at DESC чтобы самые свежие шли первыми.
# Завершённые (non-running) прогоны источника, самые свежие первыми.
rows = db.execute(
text(
"""
@ -171,31 +224,21 @@ def _alert_if_consecutive_failures(db: Session, source: str) -> None:
LIMIT :limit
"""
),
{"source": source, "limit": n + 1},
{"source": source, "limit": STREAK_SCAN_LIMIT},
).fetchall()
if len(rows) < n:
# Ещё не набралось N завершённых запусков вообще — алерт не нужен.
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):
return
# Первые N должны быть все failed/banned.
first_n = rows[:n]
if not all(r.status in ("failed", "banned") for r in first_n):
return
# (N+1)-й запуск, если есть, тоже должен НЕ быть failed/banned — иначе мы уже
# должны были отправить алерт раньше и не стоит дублировать.
if len(rows) > n and rows[n].status in ("failed", "banned"):
return
# Стрик ровно достиг порога — отправляем алерт.
sentry_sdk.capture_message(
f"Scraper source '{source}' has {n} consecutive failed/banned runs — "
f"Scraper source '{source}' has {streak} consecutive failed/banned runs — "
"manual intervention may be required (expired cookies / ban / broken parser).",
level="error",
)
logger.error(
"sentry alert sent: source=%s has %d consecutive failed/banned runs", source, n
"sentry alert sent: source=%s has %d consecutive failed/banned runs", source, streak
)
except Exception:
pass # sentry_sdk not initialised in dev, or query failed — best-effort only
@ -210,8 +253,8 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None:
Причина обычно капча/пустая выдача источника, у которого нет (или не сработал)
детект блокировки (см. providers/cian/serp.py, providers/yandex/serp.py).
Anti-spam: тот же N-й-стрик паттерн, что у _alert_if_consecutive_failures
алерт срабатывает ровно когда стрик достигает порога, не на каждом запуске сверх.
Anti-spam: та же разреженная лестница вех, что у _alert_if_consecutive_failures
(#2670) — N, 2N, 4N…, а не «ровно N и дальше тишина».
#2703: анти-спам «один раз на стрик» безопасен ТОЛЬКО там, где стрик может
прерваться. Сторож читал колонку total_seen (DEFAULT 0), которой у 28 из 53
@ -241,10 +284,10 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None:
LIMIT :limit
"""
),
{"source": source, "limit": n + 1},
{"source": source, "limit": STREAK_SCAN_LIMIT},
).fetchall()
if len(rows) < n:
if not rows:
return
def _is_zero_done(r: Any) -> bool:
@ -262,15 +305,13 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None:
_warn_source_has_no_result_metric(source, tuple(sorted(rows[0].counters or {})))
return
first_n = rows[:n]
if not all(_is_zero_done(r) for r in first_n):
return
if len(rows) > n and _is_zero_done(rows[n]):
streak = _leading_streak(rows, _is_zero_done)
capped = streak >= STREAK_SCAN_LIMIT
if not capped and not _streak_alert_due(streak, n):
return
sentry_sdk.capture_message(
f"Scraper source '{source}' has {n} consecutive 'done' runs with zero "
f"Scraper source '{source}' has {streak} consecutive 'done' runs with zero "
"lots fetched — captcha/layout-change likely undetected "
"(manual check recommended).",
level="error",
@ -278,7 +319,7 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None:
logger.error(
"sentry alert sent: source=%s has %d consecutive zero-result 'done' runs",
source,
n,
streak,
)
except Exception:
pass # sentry_sdk not initialised in dev, or query failed — best-effort only

View file

@ -0,0 +1,295 @@
"""#2670: длинная серия неудач замолкает навсегда; оборванный прогон зовётся успехом.
**1. Анти-спам, запертый в «один раз навсегда».** Оба сторожа слали алерт РОВНО на N-й
подряд неудаче и дальше молчали. Пока источники падали вперемешку с успехами, серия
рвалась и алерт взводился заново; у постоянно сломанного источника рваться нечему.
Прод 2026-08-06:
* `avito_full_load` 31 неудача подряд (20 failed + 28 banned + 11 cancelled в
истории), последний успешный прогон 03.07, то есть 34 дня без сбора и ровно один
алерт на третьей неудаче;
* `avito_full_load_exhaustive` 5 подряд;
* `domclick_city_sweep` текущий стрик 1 при 47 завершённых прогонах: условие
прерывания у этого сторожа ДОСТИЖИМО (в отличие от #2703, где оно было
недостижимо структурно) просто у сломанного источника оно не наступает.
Правка: разреженная лестница напоминаний N, 2N, 4N и не реже, чем раз в
STREAK_ALERT_MAX_PERIOD×N прогонов.
**2. Оборванный прогон.** Признак обрыва рождался в цикле `fetch_city` по ROOM_BUCKETS
(`break` на блоке, `continue` на битом бакете) и наружу не выходил: метод возвращает
голый список лотов, а pipeline писал в `counters.pages_fetched` расчётную оценку
`бакеты × страницы` не измерение. Поэтому прогон, прошедший 3 бакета из 6, был
неотличим от полного и уходил в `done`. Теперь охват считает сам скрейпер и несёт его
наружу тем же каналом, что и `blocked`.
Фальсификация: на старом коде падают тесты лестницы (сторож молчал при стрике >N) и
тест неполного охвата (прогон с лотами и неполным охватом уходил в `mark_done`).
"""
from __future__ import annotations
import os
from types import SimpleNamespace
from typing import Any
from unittest.mock import AsyncMock, MagicMock, patch
import pytest
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
from scraper_kit.orchestration import runs as kit_runs
from scraper_kit.orchestration.pipeline import run_domclick_city_sweep
from scraper_kit.providers.domclick.serp import ROOM_BUCKETS, DomClickScraper
from app.services import scrape_runs as app_runs
_MODULES = {"kit": kit_runs, "app": app_runs}
PFX = "scraper_kit.orchestration.pipeline"
def _db(rows: list[Any]) -> MagicMock:
db = MagicMock()
db.execute.return_value.fetchall.return_value = rows
return db
# ── 1. Лестница напоминаний ──────────────────────────────────────────────────
@pytest.mark.parametrize("name", list(_MODULES))
@pytest.mark.parametrize(
("streak", "due"),
[
(0, False),
(2, False),
(3, True), # первый алерт там же, где и раньше
(4, False),
(5, False),
(6, True), # 2N
(9, False),
(12, True), # 4N
(24, True), # 8N
(31, False), # прод-стрик avito_full_load — между вехами
(48, True), # 16N, дальше лестница линейная
(96, True),
(144, True),
(150, False),
],
)
def test_streak_alert_ladder(name: str, streak: int, due: bool) -> None:
"""N, 2N, 4N, 8N… и дальше каждые STREAK_ALERT_MAX_PERIOD×N — не «ровно N»."""
mod = _MODULES[name]
assert mod._streak_alert_due(streak, 3) is due
# ── 2. Сторож неудач: длинная серия не замолкает ─────────────────────────────
def _fail_rows(streak: int, tail: int = 3) -> list[SimpleNamespace]:
"""Свежие сверху: `streak` неудач подряд, затем успешные прогоны."""
return [SimpleNamespace(status="banned") for _ in range(streak)] + [
SimpleNamespace(status="done") for _ in range(tail)
]
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), "avito_full_load")
return sentry
@pytest.mark.parametrize("name", list(_MODULES))
@pytest.mark.parametrize("streak", [3, 6, 12, 24, 48])
def test_failure_watchdog_keeps_reminding(name: str, streak: int) -> None:
"""На старом коде алерт был только при streak == 3; остальные вехи молчали."""
sentry = _run_failure_watchdog(_MODULES[name], _fail_rows(streak))
sentry.capture_message.assert_called_once()
assert f"{streak} consecutive" in sentry.capture_message.call_args[0][0]
@pytest.mark.parametrize("name", list(_MODULES))
@pytest.mark.parametrize("streak", [0, 1, 2, 4, 31])
def test_failure_watchdog_silent_between_milestones(name: str, streak: int) -> None:
"""Анти-спам сохраняется: между вехами сторож молчит."""
sentry = _run_failure_watchdog(_MODULES[name], _fail_rows(streak))
sentry.capture_message.assert_not_called()
@pytest.mark.parametrize("name", list(_MODULES))
def test_failure_streak_is_broken_by_success(name: str) -> None:
"""Успешный прогон рвёт стрик — условие прерывания достижимо (прод: domclick, 1)."""
mod = _MODULES[name]
rows = [
SimpleNamespace(status="banned"),
SimpleNamespace(status="banned"),
SimpleNamespace(status="done"), # рвёт: дальше 10 неудач уже не в счёт
*[SimpleNamespace(status="failed") for _ in range(10)],
]
_run_failure_watchdog(mod, rows).capture_message.assert_not_called()
@pytest.mark.parametrize("name", list(_MODULES))
def test_failure_watchdog_alerts_when_scan_window_is_full(name: str) -> None:
"""Стрик длиннее окна сканирования — молчать нельзя, каким бы ни было число."""
mod = _MODULES[name]
rows = [SimpleNamespace(status="failed") for _ in range(mod.STREAK_SCAN_LIMIT)]
_run_failure_watchdog(mod, rows).capture_message.assert_called_once()
# ── 3. Сторож нулей: та же лестница, семантика #2703 цела ────────────────────
def _zero_rows(streak: int, tail: int = 3) -> list[SimpleNamespace]:
return [SimpleNamespace(status="done", counters={"lots_fetched": 0}) for _ in range(streak)] + [
SimpleNamespace(status="done", counters={"lots_fetched": 42}) for _ in range(tail)
]
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_full_load")
return sentry
@pytest.mark.parametrize("name", list(_MODULES))
@pytest.mark.parametrize(("streak", "called"), [(2, False), (3, True), (6, True), (7, False)])
def test_zero_watchdog_uses_the_same_ladder(name: str, streak: int, called: bool) -> None:
sentry = _run_zero_watchdog(_MODULES[name], _zero_rows(streak))
assert sentry.capture_message.called is called
@pytest.mark.parametrize("name", list(_MODULES))
def test_zero_watchdog_unmeasured_still_breaks_the_streak(name: str) -> None:
"""#2703 не отменяется: «не измерено» рвёт стрик, а не копит его."""
mod = _MODULES[name]
rows = [
SimpleNamespace(status="done", counters={"lots_fetched": 0}),
SimpleNamespace(status="done", counters={"attempted": 5}), # метрики нет
*[SimpleNamespace(status="done", counters={"lots_fetched": 0}) for _ in range(10)],
]
_run_zero_watchdog(mod, rows).capture_message.assert_not_called()
# ── 4. Охват прогона доезжает из скрейпера ───────────────────────────────────
class _FakeFetcher:
async def __aenter__(self) -> _FakeFetcher:
return self
async def __aexit__(self, *args: object) -> None:
return None
def report_ban(self, reason: str) -> None:
return None
@pytest.fixture
def _no_browser(monkeypatch: pytest.MonkeyPatch) -> None:
monkeypatch.setattr(
"scraper_kit.providers._base.build_browser_fetcher",
lambda config, source: _FakeFetcher(),
)
async def test_full_sweep_reports_full_coverage(_no_browser: None) -> None:
scraper = DomClickScraper(SimpleNamespace(browser_http_endpoint="http://x:9000"))
with patch.object(DomClickScraper, "_sweep_bucket", AsyncMock(return_value=None)):
await scraper.fetch_city(city_id=4)
assert scraper.buckets_completed == scraper.buckets_total == len(ROOM_BUCKETS)
async def test_broken_bucket_is_not_counted_as_covered(_no_browser: None) -> None:
"""Бакет, упавший на разборе, пропускается (continue) — это и есть обрыв охвата."""
scraper = DomClickScraper(SimpleNamespace(browser_http_endpoint="http://x:9000"))
calls = {"n": 0}
async def _sweep(self: DomClickScraper, **_: object) -> None:
calls["n"] += 1
if calls["n"] in (2, 5):
raise ValueError("bad BFF shape")
with patch.object(DomClickScraper, "_sweep_bucket", _sweep):
await scraper.fetch_city(city_id=4)
assert scraper.buckets_completed == len(ROOM_BUCKETS) - 2
assert scraper.fetch_errors == 2
# ── 5. Неполный охват перестаёт быть «успехом» ───────────────────────────────
class _RunsRecorder:
def __init__(self) -> None:
self.calls: list[str] = []
def is_cancelled(self, db: Any, run_id: int) -> bool:
return False
def update_heartbeat(self, db: Any, run_id: int, counters: dict[str, Any]) -> None:
self.calls.append("update_heartbeat")
def mark_done(self, db: Any, run_id: int, counters: dict[str, Any]) -> None:
self.calls.append("mark_done")
def mark_failed(self, db: Any, run_id: int, error: str, counters: dict[str, Any]) -> None:
self.calls.append("mark_failed")
def mark_banned(
self, db: Any, run_id: int, error: str, counters: dict[str, Any], **kw: Any
) -> None:
self.calls.append("mark_banned")
async def _drive(*, lots_n: int, done: int, total: int, blocked: bool = False) -> list[str]:
recorder = _RunsRecorder()
lots = [MagicMock() for _ in range(lots_n)]
scraper = MagicMock()
scraper.__aenter__ = AsyncMock(return_value=scraper)
scraper.__aexit__ = AsyncMock(return_value=None)
scraper.fetch_city = AsyncMock(return_value=lots)
scraper.blocked = blocked
scraper.geo_filtered = 0
scraper.fetch_errors = total - done
scraper.buckets_completed = done
scraper.buckets_total = total
with (
patch(f"{PFX}.DomClickScraper", return_value=scraper),
patch(f"{PFX}.save_listings", MagicMock(return_value=(lots_n, 0))),
patch(f"{PFX}.runs", recorder),
):
await run_domclick_city_sweep(
MagicMock(),
config=SimpleNamespace(browser_http_endpoint="http://x:9000"),
matcher=MagicMock(),
run_id=1,
city_id=4,
pages=1,
request_delay_sec=0.0,
)
return recorder.calls
async def test_partial_coverage_with_lots_is_not_done() -> None:
"""Прогон, прошедший 3 бакета из 6, не «успешен», даже если лоты есть.
На старом коде эта ветка отсутствовала и прогон уходил в mark_done.
"""
assert (await _drive(lots_n=340, done=3, total=6))[-1] == "mark_failed"
async def test_full_coverage_with_lots_stays_done() -> None:
"""Анти-оверрич: полный охват — по-прежнему done."""
assert (await _drive(lots_n=340, done=6, total=6))[-1] == "mark_done"
async def test_block_still_wins_over_coverage() -> None:
"""Блок проверяется раньше охвата: диагноз «нас прервали снаружи» точнее (#2657)."""
assert (await _drive(lots_n=39, done=2, total=6, blocked=True))[-1] == "mark_banned"
async def test_unknown_coverage_does_not_invent_a_verdict() -> None:
"""Скрейпер не успел создаться (0/0) — судить об охвате нечем, ветка не срабатывает."""
assert (await _drive(lots_n=12, done=0, total=0))[-1] == "mark_done"

View file

@ -39,6 +39,7 @@ from scraper_kit.orchestration.pipeline import (
run_yandex_city_sweep,
run_yandex_full_load,
)
from scraper_kit.providers.domclick.serp import ROOM_BUCKETS
PFX = "scraper_kit.orchestration.pipeline"
@ -413,6 +414,10 @@ async def _drive_domclick(
blocked=blocked,
geo_filtered=0,
fetch_errors=fetch_errors,
# #2670: полный охват по умолчанию — эти фикстуры про блок/ошибки, не про обрыв
# (частичный охват проверяется в test_2670_streak_and_partial_coverage.py).
buckets_completed=len(ROOM_BUCKETS),
buckets_total=len(ROOM_BUCKETS),
)
save_mock = MagicMock(side_effect=[(lots_n, 0)] if lots_n else [])
if capture is not None:

View file

@ -3705,6 +3705,10 @@ class DomClickCitySweepCounters:
errors_count: int = 0
blocked: int = 0 # 1 если QRATOR-блок был во время sweep
geo_filtered: int = 0 # число офферов отфильтрованных geo-guard
# #2670: измеренный охват прогона (сколько комнатных бакетов пройдено из скольких).
# 0/0 = скрейпер не успел создаться — тогда судить об охвате нечем.
buckets_completed: int = 0
buckets_total: int = 0
def to_dict(self) -> dict[str, int]:
return {f.name: getattr(self, f.name) for f in fields(self)}
@ -3713,8 +3717,9 @@ class DomClickCitySweepCounters:
# Дефолтные параметры sweep'а (EKB city_id=4).
DOMCLICK_DEFAULT_CITY_ID: int = 4
DOMCLICK_DEFAULT_ROOMS: list[int] = [0, 1, 2, 3, 4] # vestigial; scraper sweeps all buckets
# Число BFF-бакетов (st/1/2/3/4/5+) — фиксировано; используется для watchdog.
_DOMCLICK_NUM_BUCKETS: int = 6
# Число BFF-бакетов (st/1/2/3/4/5+) — используется для watchdog. Берётся из самого
# ROOM_BUCKETS: два литерала, обязанных совпадать, однажды уже разъезжались (#2674).
_DOMCLICK_NUM_BUCKETS: int = len(ROOM_BUCKETS)
# Оценка времени одного fetch'а (network + parse) для watchdog.
_DOMCLICK_PER_FETCH_S: float = 12.0
# Буфер сверху расчётного бюджета (cold browser start, save-фаза, price-splits).
@ -3836,6 +3841,10 @@ async def run_domclick_city_sweep(
counters.geo_filtered = _s.geo_filtered
# Не-block fetch-ошибки скрейпера учитываем в errors_count.
counters.errors_count += _s.fetch_errors
# #2670: охват читаем даже когда фаза была снята по таймауту — ссылка на
# скрейпер живая, а его счётчик показывает, докуда прогон дошёл.
counters.buckets_completed = _s.buckets_completed
counters.buckets_total = _s.buckets_total
# pages_fetched: worst-case число страниц (buckets × pages cap).
counters.pages_fetched = _num_fetches
@ -3868,6 +3877,29 @@ async def run_domclick_city_sweep(
"collected before abort (#2657)",
counters.to_dict(),
)
elif 0 < counters.buckets_completed < counters.buckets_total:
# #2670: прогон оборван на середине — прошёл часть комнатных бакетов и
# бросил остальные (битый бакет → continue, снятие фазы по таймауту).
# Лоты у него есть, известной ошибки нет — и до этой ветки он отчитывался
# успехом. «Успех» определялся как «не поймали известную ошибку», а не как
# «сделали то, что собирались»: сравниваем с ожидаемым охватом, не с нулём.
logger.error(
"domclick-sweep run_id=%d: пройдено %d бакетов из %d "
"(lots=%d, errors=%d) — marking failed (#2670)",
run_id,
counters.buckets_completed,
counters.buckets_total,
counters.lots_fetched,
counters.errors_count,
)
runs.mark_failed(
db,
run_id,
f"sweep оборван: пройдено {counters.buckets_completed} комнатных "
f"бакетов из {counters.buckets_total}, собрано "
f"{counters.lots_fetched} лотов (#2670)",
counters.to_dict(),
)
elif counters.lots_fetched == 0 and counters.errors_count > 0:
logger.error(
"domclick-sweep run_id=%d: 0 listings with errors=%d — marking failed",
@ -3885,11 +3917,13 @@ async def run_domclick_city_sweep(
logger.info(
"domclick-sweep run_id=%d done: lots=%d (ins=%d/upd=%d) "
"pages=%d errors=%d blocked=%d geo_filtered=%d",
"buckets=%d/%d pages=%d errors=%d blocked=%d geo_filtered=%d",
run_id,
counters.lots_fetched,
counters.lots_inserted,
counters.lots_updated,
counters.buckets_completed,
counters.buckets_total,
counters.pages_fetched,
counters.errors_count,
counters.blocked,

View file

@ -50,6 +50,52 @@ CONSECUTIVE_FAILURE_ALERT_THRESHOLD = 3
# невидим для CONSECUTIVE_FAILURE_ALERT_THRESHOLD (тот считает только failed/banned).
CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD = 3
# #2670: анти-спам «один раз на стрик» безопасен ТОЛЬКО там, где стрик прерывается
# не только в принципе, но и на практике. Оба сторожа ниже слали алерт ровно на N-й
# подряд неудаче и дальше молчали навсегда — а у постоянно сломанного источника
# «дальше» длится месяцами. Прод 2026-08-06: у avito_full_load 31 неудача подряд,
# последний успешный прогон 03.07 (34 дня без сбора), алерт был ровно один — на
# третьей; у avito_full_load_exhaustive 5 подряд. Тишина при этом неотличима от
# «всё хорошо» — ровно та ловушка, из-за которой #2574 месяц выглядела как норма.
#
# Вместо «ровно N» — разреженная лестница напоминаний: N, 2N, 4N, 8N…, а дальше не
# реже, чем раз в STREAK_ALERT_MAX_PERIOD×N прогонов. Лестница по ПРОГОНАМ, а не
# «раз в сутки», потому что источники идут разным тактом: domclick_city_sweep — раз
# в день, proxy_healthcheck — раз в полчаса; календарное разрежение для одного из
# них всегда будет либо спамом, либо молчанием.
STREAK_ALERT_MAX_PERIOD = 16
# Потолок сканирования истории источника при подсчёте стрика. Достигнутый потолок
# сам по себе повод для алерта (стрик заведомо огромен) — так «замолчать навсегда»
# невозможно по построению, а не по счастливому совпадению чисел.
STREAK_SCAN_LIMIT = 500
def _streak_alert_due(streak: int, threshold: int) -> bool:
"""Достиг ли стрик очередной вехи напоминания (#2670).
True на threshold, 2×, 4×, 8× и дальше на каждом кратном
STREAK_ALERT_MAX_PERIOD×threshold. Первый алерт приходит там же, где и раньше
на N-й подряд неудаче; меняется только то, что он не последний.
"""
if streak < threshold or streak % threshold:
return False
mult = streak // threshold
if mult % STREAK_ALERT_MAX_PERIOD == 0:
return True
return mult & (mult - 1) == 0
def _leading_streak(rows: list[Any], is_bad: Callable[[Any], bool]) -> int:
"""Длина серии подряд идущих «плохих» строк с начала списка (свежие — первыми)."""
streak = 0
for row in rows:
if not is_bad(row):
break
streak += 1
return streak
# #2686: диагноз оборванного прогона. Пишется в scrape_runs.ban_kind (миграция 218)
# РЯДОМ со status='banned', а не ВМЕСТО него — сознательный выбор между «новый
# статус» и «явное поле причины»:
@ -142,12 +188,18 @@ def _column_counts(counters: dict[str, int]) -> tuple[int | None, int | None]:
def _alert_if_consecutive_failures(db: Session, source: str) -> None:
"""Отправить Sentry alert если последние CONSECUTIVE_FAILURE_ALERT_THRESHOLD
завершённых запусков для данного source имеют статус 'failed' или 'banned'.
"""Sentry alert на серию из CONSECUTIVE_FAILURE_ALERT_THRESHOLD неудач подряд
(статусы 'failed'/'banned') у данного source.
Anti-spam: алерт срабатывает ТОЛЬКО когда стрик РОВНО равен порогу т.е. запрос
возвращает ровно N последних (failed|banned) и (N+1)-й, если существует, НЕ является
failed/banned. Это предотвращает повторный алерт на каждой ошибке сверх порога.
Anti-spam: не на каждой неудаче, а по разреженной лестнице вех (см.
_streak_alert_due). До #2670 алерт приходил РОВНО на N-й неудаче и дальше не
повторялся никогда: серия, ставшая длиннее порога, замолкала навсегда. На проде
это дало avito_full_load 31 неудача подряд, 34 дня без единого успешного
прогона, один алерт за всё время.
Стрик прерывается любым завершением, кроме failed/banned, по данным прода это
достижимо и достигается (у domclick_city_sweep текущий стрик равен 1 при 47
завершённых прогонах), поэтому лестница не вырождается в постоянный алерт.
Best-effort: весь блок обёрнут в try/except сбой запроса или неинициализированный
Sentry НЕ должен нарушать вызывающий mark_* путь.
@ -156,8 +208,7 @@ def _alert_if_consecutive_failures(db: Session, source: str) -> None:
return
n = CONSECUTIVE_FAILURE_ALERT_THRESHOLD
try:
# Берём последние N+1 завершённых (non-running) запусков по source.
# Сортируем по finished_at DESC чтобы самые свежие шли первыми.
# Завершённые (non-running) прогоны источника, самые свежие первыми.
rows = db.execute(
text(
"""
@ -168,31 +219,21 @@ def _alert_if_consecutive_failures(db: Session, source: str) -> None:
LIMIT :limit
"""
),
{"source": source, "limit": n + 1},
{"source": source, "limit": STREAK_SCAN_LIMIT},
).fetchall()
if len(rows) < n:
# Ещё не набралось N завершённых запусков вообще — алерт не нужен.
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):
return
# Первые N должны быть все failed/banned.
first_n = rows[:n]
if not all(r.status in ("failed", "banned") for r in first_n):
return
# (N+1)-й запуск, если есть, тоже должен НЕ быть failed/banned — иначе мы уже
# должны были отправить алерт раньше и не стоит дублировать.
if len(rows) > n and rows[n].status in ("failed", "banned"):
return
# Стрик ровно достиг порога — отправляем алерт.
sentry_sdk.capture_message(
f"Scraper source '{source}' has {n} consecutive failed/banned runs — "
f"Scraper source '{source}' has {streak} consecutive failed/banned runs — "
"manual intervention may be required (expired cookies / ban / broken parser).",
level="error",
)
logger.error(
"sentry alert sent: source=%s has %d consecutive failed/banned runs", source, n
"sentry alert sent: source=%s has %d consecutive failed/banned runs", source, streak
)
except Exception:
pass # sentry_sdk not initialised in dev, or query failed — best-effort only
@ -207,8 +248,8 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None:
Причина обычно капча/пустая выдача источника, у которого нет (или не сработал)
детект блокировки (см. providers/cian/serp.py, providers/yandex/serp.py).
Anti-spam: тот же N-й-стрик паттерн, что у _alert_if_consecutive_failures
алерт срабатывает ровно когда стрик достигает порога, не на каждом запуске сверх.
Anti-spam: та же разреженная лестница вех, что у _alert_if_consecutive_failures
(#2670) — N, 2N, 4N…, а не «ровно N и дальше тишина».
#2703: анти-спам «один раз на стрик» безопасен ТОЛЬКО там, где стрик может
прерваться. Сторож читал колонку total_seen (DEFAULT 0), которой у 28 из 53
@ -240,10 +281,10 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None:
LIMIT :limit
"""
),
{"source": source, "limit": n + 1},
{"source": source, "limit": STREAK_SCAN_LIMIT},
).fetchall()
if len(rows) < n:
if not rows:
return
def _is_zero_done(r: Any) -> bool:
@ -261,15 +302,13 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None:
_warn_source_has_no_result_metric(source, tuple(sorted(rows[0].counters or {})))
return
first_n = rows[:n]
if not all(_is_zero_done(r) for r in first_n):
return
if len(rows) > n and _is_zero_done(rows[n]):
streak = _leading_streak(rows, _is_zero_done)
capped = streak >= STREAK_SCAN_LIMIT
if not capped and not _streak_alert_due(streak, n):
return
sentry_sdk.capture_message(
f"Scraper source '{source}' has {n} consecutive 'done' runs with zero "
f"Scraper source '{source}' has {streak} consecutive 'done' runs with zero "
"lots fetched — captcha/layout-change likely undetected "
"(manual check recommended).",
level="error",
@ -277,7 +316,7 @@ def _alert_if_consecutive_zero_results(db: Session, source: str) -> None:
logger.error(
"sentry alert sent: source=%s has %d consecutive zero-result 'done' runs",
source,
n,
streak,
)
except Exception:
pass # sentry_sdk not initialised in dev, or query failed — best-effort only

View file

@ -235,6 +235,9 @@ class DomClickScraper(BaseScraper):
geo_filtered офферы вне ЕКБ bbox или с неверным offerRegionName
blocked True если sweep был прерван QRATOR-блоком
fetch_errors не-block ошибки извлечения JSON (truncated/garbled/bad shape)
buckets_total сколько комнатных бакетов прогон собирался пройти
buckets_completed сколько прошёл ФАКТИЧЕСКИ (#2670); меньше total = прогон
оборван, сколько бы лотов он ни успел взять
"""
name = "domklik"
@ -263,6 +266,14 @@ class DomClickScraper(BaseScraper):
# структура). В отличие от parse_failures (per-item), это per-fetch ошибки,
# которые ограничивают сбор бакета. Учитываются в honest-status pipeline.
self.fetch_errors: int = 0
# #2670: охват прогона. `blocked` отвечает на «нас прервали снаружи», эти два —
# на «сколько работы прогон реально сделал». Признак обрыва рождается ЗДЕСЬ, в
# цикле по ROOM_BUCKETS (break на блоке, continue на битом бакете), и до #2670
# наружу не выходил: fetch_city возвращает голый список лотов, а pipeline писал
# в counters.pages_fetched расчётную оценку buckets × pages — не измерение.
# Поэтому прогон, прошедший 3 бакета из 6, был неотличим от полного.
self.buckets_total: int = len(ROOM_BUCKETS)
self.buckets_completed: int = 0
async def __aenter__(self) -> DomClickScraper:
await super().__aenter__()
@ -356,12 +367,15 @@ class DomClickScraper(BaseScraper):
exc_info=True,
)
continue
self.buckets_completed += 1
logger.info(
"domklik: fetch_city done city_id=%d total=%d "
"domklik: fetch_city done city_id=%d total=%d buckets=%d/%d "
"parse_failures=%d geo_filtered=%d fetch_errors=%d blocked=%s",
city_id,
len(out_lots),
self.buckets_completed,
self.buckets_total,
self.parse_failures,
self.geo_filtered,
self.fetch_errors,