fix(tradein/scrapers): обрыв по серии блоков рвал каждый прогон, включая здоровые #3188
4 changed files with 516 additions and 32 deletions
|
|
@ -1114,6 +1114,32 @@ class Settings(BaseSettings):
|
|||
default=90.0, validation_alias="AVITO_DETAIL_FETCH_TIMEOUT_S"
|
||||
)
|
||||
|
||||
# #3184: доля блоков в скользящем окне последних N попыток -- критерий обрыва
|
||||
# avito_detail_backfill (app.services.backfill_block_breaker.BlockRatioBreaker),
|
||||
# взамен голого "N блоков подряд". ТОЛЬКО avito -- изначальный план распространить
|
||||
# тот же критерий на domclick_detail_backfill снят ревью (#3184 review MAJOR 1):
|
||||
# у Домклика один выделенный residential-прокси БЕЗ ротации (см.
|
||||
# data/sql/175_scrape_schedules_seed_domclick_detail_backfill.sql), калибровка
|
||||
# 20/0.7 сделана на пуле С ротацией (avito) и для него не годится без своего
|
||||
# замера -- отдельная задача.
|
||||
# Калибровка по 40 прогонам avito за 14 суток (2026-08-14..28): ВСЕ 40
|
||||
# закончились 'banned', доля блоков колебалась 25-100% (25%/190 попыток честно
|
||||
# обогащённых 125, 39%/100, 46%/71, 32%/63, 100%/5) -- "N подряд" не отличал
|
||||
# выгоревший прокси-пул от здорового прогона, потому что блоки автокоррелированы
|
||||
# и пачки 5-6 подряд встречаются в каждом прогоне при базовой доле 25-46%. Окно
|
||||
# 20 выбрано заметно больше типичной пачки, порог 0.7 -- между здоровыми (25-48%)
|
||||
# и выгоревшими (100%) прогонами большой запас. ENV: DETAIL_BACKFILL_BLOCK_RATIO_
|
||||
# WINDOW / _THRESHOLD -- подкрутка без релиза.
|
||||
# ge/le (#3184 review MINOR 2): WINDOW=0 или THRESHOLD в процентах (70 вместо
|
||||
# 0.7, частая опечатка оператора) иначе молча выключают критерий навсегда --
|
||||
# 0.7 > 70 никогда не бывает True, а окно 0 схлопывает should_abort() в no-op.
|
||||
detail_backfill_block_ratio_window: int = Field(
|
||||
default=20, ge=1, validation_alias="DETAIL_BACKFILL_BLOCK_RATIO_WINDOW"
|
||||
)
|
||||
detail_backfill_block_ratio_threshold: float = Field(
|
||||
default=0.7, ge=0.0, le=1.0, validation_alias="DETAIL_BACKFILL_BLOCK_RATIO_THRESHOLD"
|
||||
)
|
||||
|
||||
# ── #884/#905/#1805: BrowserFetcher — HTTP-клиент к tradein-browser ─────────
|
||||
# scraper_fetch_mode: "browser" (дефолт с #1805 — HTTP POST к tradein-browser
|
||||
# /fetch через per-provider camoufox + ротирующий backconnect-прокси) или
|
||||
|
|
|
|||
133
tradein-mvp/backend/app/services/backfill_block_breaker.py
Normal file
133
tradein-mvp/backend/app/services/backfill_block_breaker.py
Normal file
|
|
@ -0,0 +1,133 @@
|
|||
"""Критерий обрыва avito_detail_backfill по доле блоков (#3184).
|
||||
|
||||
Модуль написан источник-агностично (см. BlockRatioBreaker ниже), но единственный
|
||||
текущий вызывающий -- avito_detail_backfill. Исходная задача планировала тот же
|
||||
критерий и для domclick_detail_backfill ("тот же критерий, что у avito"), но это
|
||||
снято ревью (#3184 review MAJOR 1): у Домклика ОДИН выделенный residential-прокси
|
||||
БЕЗ ротации (data/sql/175_scrape_schedules_seed_domclick_detail_backfill.sql,
|
||||
max_consecutive_blocks=3 против avito=5 -- намеренно ниже, ранний обрыв бережёт
|
||||
репутацию единственного узла), а калибровка окна/порога ниже сделана на avito-пуле
|
||||
С ротацией. При устойчивой доле блоков около порога (например 0.65) Домклик по
|
||||
этому критерию не обрывался бы никогда -- десятки запросов по QRATOR с одного
|
||||
IP, ровно то, что max_consecutive_blocks=3 был призван предотвратить. Домклик
|
||||
получит свою калибровку отдельной задачей.
|
||||
|
||||
Раньше avito_detail_backfill абортил по голой длине "N блоков подряд"
|
||||
(max_consecutive_blocks). Замер 40 прогонов avito за 14 суток (2026-08-14..28)
|
||||
показал: ВСЕ 40 закончились статусом 'banned', при этом доля блоков в них
|
||||
колебалась 25-100% (47/190=25%, 39/100=39%, 33/71=46%, 20/63=32%, 5/5=100%) --
|
||||
"N подряд" не отличал выгоревший прокси-пул от прогона, который честно
|
||||
обогатил 125 карточек из 190. Причина: блоки автокоррелированы и идут
|
||||
пачками -- при базовой доле 25-46% пачка из 5-6 подряд у независимой модели
|
||||
была бы редкостью, но встречается в каждом прогоне.
|
||||
|
||||
Критерий заменён на долю блоков в скользящем окне последних N попыток: пачка
|
||||
сама по себе больше не абортит, абортит устойчиво высокая доля (>= порога,
|
||||
не строго "выше" -- окно 20 при пороге 0.7 абортит уже на 14/20, 13/20 ещё
|
||||
нет). Окно короче типичной пачки бесполезно, поэтому дефолт (settings.
|
||||
detail_backfill_block_ratio_window) заметно больше 5-6.
|
||||
|
||||
Safety-net для прогонов КОРОЧЕ окна (снапшот меньше window_size -- ratio-
|
||||
критерий физически недостижим, окно никогда не заполнится) сохранён: если ВСЕ
|
||||
попытки с начала прогона были блоками (ни одного успеха) и накопилось не
|
||||
меньше safety_min -- абортим не дожидаясь окна. Раньше это было единственным
|
||||
критерием под именем max_consecutive_blocks; имя и дефолт (5) сохранены как
|
||||
safety_min, чтобы не менять поведение коротких прогонов вида 0/5 (#3184).
|
||||
|
||||
Условие "снапшот короче окна" обязательно, иначе safety-net гасит сам фикс:
|
||||
холодный старт сессии/прокси -- самое вероятное место пачки блоков, и на
|
||||
длинном прогоне (снапшот >= окна) пачка из safety_min блоков В НАЧАЛЕ
|
||||
абортила бы точно так же, как до правки, до того как ratio-критерий вообще
|
||||
успел бы включиться (#3184 review MAJOR 2).
|
||||
|
||||
Отказы-не-блоки (сеть/таймаут/парсер уровня "ответ есть, но наш") НЕ рвут
|
||||
"чистоту" серии для safety-net -- это исторический инвариант
|
||||
avito_detail_backfill (только успех сбрасывал consecutive_blocks,
|
||||
TimeoutError/Exception его не трогали), сохранён через record_failure():
|
||||
она двигает окно (значит, доля в ratio-критерии всё равно снижается), но НЕ
|
||||
трогает _pure_block_run и НЕ сбрасывает серию. Из-за этого 4 таймаута + 5
|
||||
блоков подряд (ни одного успеха) всё ещё дают safety-net abort -- это
|
||||
осознанное поведение, а не баг: таймаут не доказывает, что тракт живой ответ
|
||||
получает, только успех доказывает.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections import Counter, deque
|
||||
from dataclasses import dataclass, field
|
||||
|
||||
|
||||
@dataclass
|
||||
class BlockRatioBreaker:
|
||||
"""Стейт одного прогона detail-backfill'а: окно, safety-net, гистограмма пачек."""
|
||||
|
||||
window_size: int
|
||||
ratio_threshold: float
|
||||
safety_min: int
|
||||
snapshot_size: int
|
||||
_window: deque[bool] = field(init=False, repr=False)
|
||||
_consecutive_blocks: int = field(default=0, init=False)
|
||||
_pure_block_run: bool = field(default=True, init=False)
|
||||
streak_histogram: Counter[int] = field(default_factory=Counter, init=False)
|
||||
|
||||
def __post_init__(self) -> None:
|
||||
self._window = deque(maxlen=max(self.window_size, 1))
|
||||
|
||||
@property
|
||||
def consecutive_blocks(self) -> int:
|
||||
"""Для логов ABORT/BLOCKED -- та же цифра, что раньше выводилась в consecutive=%d."""
|
||||
return self._consecutive_blocks
|
||||
|
||||
def record_block(self) -> None:
|
||||
self._consecutive_blocks += 1
|
||||
self._window.append(True)
|
||||
|
||||
def record_success(self) -> None:
|
||||
"""Единственный исход, снимающий safety-net (#3184: пачка блоков ПОСЛЕ хотя
|
||||
бы одного успеха -- уже не "чистый с рождения прогона" burst)."""
|
||||
self._flush_streak()
|
||||
self._pure_block_run = False
|
||||
self._window.append(False)
|
||||
|
||||
def record_failure(self) -> None:
|
||||
"""Отказ-не-блок: в знаменатель окна идёт (тракт получил ответ, не блок),
|
||||
но серию/purity НЕ трогает -- зеркалит исторический avito-инвариант."""
|
||||
self._window.append(False)
|
||||
|
||||
def record_neutral(self) -> None:
|
||||
"""404/gone у avito -- не блок и не проверка тракта, не двигает ни окно, ни
|
||||
серию (тот же смысл, что "neutral to the breaker" в комментарии у
|
||||
AvitoListingGoneError). Класс написан источник-агностично на будущее (не
|
||||
только avito), но сейчас единственный вызывающий -- avito_detail_backfill
|
||||
(#3184 review MAJOR 1: применение к domclick_detail_backfill снято из этой
|
||||
задачи -- своя калибровка, свои ограничения прокси-пула)."""
|
||||
return
|
||||
|
||||
def _flush_streak(self) -> None:
|
||||
if self._consecutive_blocks:
|
||||
self.streak_histogram[self._consecutive_blocks] += 1
|
||||
self._consecutive_blocks = 0
|
||||
|
||||
def should_abort(self) -> bool:
|
||||
# Safety-net -- ТОЛЬКО когда ratio-критерий физически недостижим (снапшот
|
||||
# короче окна), иначе пачка safety_min блоков в начале длинного прогона
|
||||
# абортила бы его так же, как до правки (#3184 review MAJOR 2).
|
||||
if (
|
||||
self.snapshot_size < self.window_size
|
||||
and self._pure_block_run
|
||||
and self._consecutive_blocks >= self.safety_min
|
||||
):
|
||||
return True
|
||||
if len(self._window) == self.window_size and self.window_size > 0:
|
||||
ratio = sum(self._window) / self.window_size
|
||||
if ratio >= self.ratio_threshold:
|
||||
return True
|
||||
return False
|
||||
|
||||
def finalize(self) -> dict[str, int]:
|
||||
"""Досчитать хвостовую пачку (прогон оборвался посреди серии блоков, без
|
||||
финального успеха) и отдать гистограмму для counters (jsonb) -- ключи СТРОКОЙ:
|
||||
jsonb/JSON всё равно хранит только строковые ключи, отдаём их такими сразу,
|
||||
чтобы то, что пишем ({6: 1}), совпадало с тем, что потом читаем ({"6": 1})."""
|
||||
self._flush_streak()
|
||||
return {str(streak): count for streak, count in self.streak_histogram.items()}
|
||||
|
|
@ -26,6 +26,17 @@ rotate IP on every block, abort after max_consecutive_blocks. Статус об
|
|||
3306 (blocked=5, 0 обогащено) получил 'platform' по УМОЛЧАНИЮ — финализатор
|
||||
диагноз не передавал, а browser-режим fetch_detail всё равно превращал отказ
|
||||
сайдкара в AvitoBlockedError, так что передавать было бы нечего.
|
||||
|
||||
КРИТЕРИЙ ОБРЫВА (#3184): голое "N блоков подряд" (max_consecutive_blocks)
|
||||
абортило прогон 125/190 (25% блоков, честно обогащённых 125 карточек) статусом
|
||||
'banned' наравне с прогоном 0/5 (100% блоков) — блоки идут автокоррелированными
|
||||
пачками, и пачка 5-6 подряд у здорового прогона с базовой долей 25-46% не
|
||||
редкость. Теперь абортит app.services.backfill_block_breaker.BlockRatioBreaker:
|
||||
доля блоков в скользящем окне последних N попыток (settings.detail_backfill_
|
||||
block_ratio_window/_threshold) ИЛИ safety-net для коротких прогонов (max_
|
||||
consecutive_blocks блоков подряд БЕЗ единого успеха с самого начала — старое
|
||||
поведение для случая вида 0/5). Гистограмма длин пачек блоков едет в
|
||||
counters["block_streak_histogram"], иначе эффект правки нечем измерить постфактум.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
|
@ -70,6 +81,7 @@ from sqlalchemy.orm import Session
|
|||
from app.core.config import settings
|
||||
from app.core.shutdown import shutdown_requested
|
||||
from app.services import scrape_runs as runs_mod
|
||||
from app.services.backfill_block_breaker import BlockRatioBreaker
|
||||
from app.services.proxy_egress import resolve_proxy_url
|
||||
from app.services.scraper_adapters import RealProxyProvider, RealScraperConfig
|
||||
|
||||
|
|
@ -193,7 +205,11 @@ async def run_avito_detail_backfill(
|
|||
всплеску свежих oblast-листингов вытеснить ЕКБ из top-N по scraped_at.
|
||||
budget_sec: float -- wall-clock budget per run, default 3600s.
|
||||
request_delay_sec: float -- delay between listings, default 6.0s.
|
||||
max_consecutive_blocks: int -- abort threshold, default 5.
|
||||
max_consecutive_blocks: int -- safety_min для BlockRatioBreaker (#3184):
|
||||
абортит прогоны короче окна, если ВСЕ попытки с начала были блоками
|
||||
(ни одного успеха), default 5. Основной критерий -- доля блоков в
|
||||
скользящем окне, см. settings.detail_backfill_block_ratio_window/
|
||||
_threshold (module docstring).
|
||||
max_consecutive_failures: int -- порог обрыва по отказам-не-блокам,
|
||||
default 25 (см. комментарий у чтения параметра ниже).
|
||||
|
||||
|
|
@ -205,6 +221,8 @@ async def run_avito_detail_backfill(
|
|||
oblast_batch_size = int(params.get("oblast_batch_size", 100))
|
||||
budget_sec = float(params.get("budget_sec", 3600))
|
||||
request_delay_sec = float(params.get("request_delay_sec", 6.0))
|
||||
# #3184: теперь safety_min BlockRatioBreaker (см. module docstring) -- не общий
|
||||
# критерий обрыва, а страховка для прогонов короче окна.
|
||||
max_consecutive_blocks = int(params.get("max_consecutive_blocks", 5))
|
||||
# Брейкер на отказы-НЕ-блоки. Блоки свой брейкер имели с самого начала, отказы —
|
||||
# нет, и это стоило трёх ночей подряд: 3-5 августа прогон делал ~1600 попыток,
|
||||
|
|
@ -399,7 +417,18 @@ async def run_avito_detail_backfill(
|
|||
"curl/backconnect" if use_curl else "browser",
|
||||
)
|
||||
|
||||
consecutive_blocks = 0
|
||||
# #3184: доля блоков в скользящем окне вместо голого "N подряд" -- см. module
|
||||
# docstring и app.services.backfill_block_breaker.
|
||||
breaker = BlockRatioBreaker(
|
||||
window_size=int(settings.detail_backfill_block_ratio_window),
|
||||
ratio_threshold=float(settings.detail_backfill_block_ratio_threshold),
|
||||
safety_min=max_consecutive_blocks,
|
||||
# snapshot_size гейтит safety-net (#3184 review MAJOR 2): пачка блоков в
|
||||
# начале ДЛИННОГО прогона не должна абортить его так же, как раньше --
|
||||
# safety-net включён только когда снапшот короче окна и ratio-критерий
|
||||
# физически недостижим (см. backfill_block_breaker module docstring).
|
||||
snapshot_size=len(snapshot),
|
||||
)
|
||||
consecutive_failures = 0
|
||||
aborted_by_blocks = False
|
||||
do_sleep = False
|
||||
|
|
@ -530,7 +559,7 @@ async def run_avito_detail_backfill(
|
|||
counters.enriched += 1
|
||||
if use_curl:
|
||||
items_since_warm += 1
|
||||
consecutive_blocks = 0
|
||||
breaker.record_success()
|
||||
consecutive_failures = 0
|
||||
|
||||
except AvitoListingGoneError as gone_exc:
|
||||
|
|
@ -539,10 +568,11 @@ async def run_avito_detail_backfill(
|
|||
# browser-mode рендерит их «Ошибка 404» без item-view → раньше это
|
||||
# ловилось как soft-block → consecutive-block breaker абортил run до
|
||||
# live-листингов (run 458: attempted=5 enriched=0 blocked=5 → abort).
|
||||
# 404 нейтрален к breaker'у: consecutive_blocks НЕ трогаем (не растим
|
||||
# и не сбрасываем). Метим is_active=FALSE → листинг уходит из scope
|
||||
# (snapshot SELECT фильтрует is_active = TRUE) и не тратит фетчи впредь.
|
||||
# 404 нейтрален к breaker'у (#3184: record_neutral -- ни окно, ни серия,
|
||||
# ни safety-net не трогаются). Метим is_active=FALSE → листинг уходит из
|
||||
# scope (snapshot SELECT фильтрует is_active = TRUE) и не тратит фетчи впредь.
|
||||
counters.gone += 1
|
||||
breaker.record_neutral()
|
||||
# 404 — честный ответ площадки, значит тракт цел: серия отказов
|
||||
# прерывается (блочный брейкер 404 не трогает — см. #2034).
|
||||
consecutive_failures = 0
|
||||
|
|
@ -586,7 +616,7 @@ async def run_avito_detail_backfill(
|
|||
)
|
||||
|
||||
except (AvitoBlockedError, AvitoRateLimitedError) as e:
|
||||
consecutive_blocks += 1
|
||||
breaker.record_block()
|
||||
counters.blocked += 1
|
||||
failure_census[_failure_signature(e)] += 1
|
||||
block_ban_kinds[ban_kind_of_exception(e)] += 1
|
||||
|
|
@ -596,17 +626,18 @@ async def run_avito_detail_backfill(
|
|||
run_id,
|
||||
idx + 1,
|
||||
len(snapshot),
|
||||
consecutive_blocks,
|
||||
breaker.consecutive_blocks,
|
||||
e,
|
||||
)
|
||||
# abort FIRST — на последнем блоке не тратим cooldown+research/rotate
|
||||
# (consecutive_blocks уже инкрементнут выше).
|
||||
if consecutive_blocks >= max_consecutive_blocks:
|
||||
# abort FIRST — на последнем блоке не тратим cooldown+research/rotate.
|
||||
# #3184: доля в скользящем окне ИЛИ safety-net (все попытки — блоки, ни
|
||||
# одного успеха, накопилось max_consecutive_blocks) -- см. module docstring.
|
||||
if breaker.should_abort():
|
||||
logger.error(
|
||||
"avito_detail_backfill: run_id=%d ABORT -- %d consecutive blocks, "
|
||||
"частая причина: %s. enriched=%d attempted=%d",
|
||||
run_id,
|
||||
consecutive_blocks,
|
||||
breaker.consecutive_blocks,
|
||||
_top_failure(failure_census) or "причина не определена",
|
||||
counters.enriched,
|
||||
counters.attempted,
|
||||
|
|
@ -655,6 +686,7 @@ async def run_avito_detail_backfill(
|
|||
# run не zombie #1950). Не считаем soft-блоком: rotate не дёргаем.
|
||||
counters.failed += 1
|
||||
consecutive_failures += 1
|
||||
breaker.record_failure()
|
||||
failure_census[_failure_signature(e)] += 1
|
||||
logger.warning(
|
||||
"avito_detail_backfill: run_id=%d listing %s TIMEOUT (>%.0fs) -- skip",
|
||||
|
|
@ -670,6 +702,7 @@ async def run_avito_detail_backfill(
|
|||
except Exception as e:
|
||||
counters.failed += 1
|
||||
consecutive_failures += 1
|
||||
breaker.record_failure()
|
||||
failure_census[_failure_signature(e)] += 1
|
||||
logger.warning(
|
||||
"avito_detail_backfill: run_id=%d listing %s failed: %s",
|
||||
|
|
@ -700,6 +733,12 @@ async def run_avito_detail_backfill(
|
|||
|
||||
counters.duration_sec = time.monotonic() - start
|
||||
current_counters = counters.to_dict()
|
||||
# #3184: гистограмма длин пачек блоков (streak -> сколько раз встретилась) --
|
||||
# иначе эффект правки на #2674-статистике нечем измерить постфактум. finalize()
|
||||
# досчитывает хвостовую пачку, если прогон оборвался посреди серии.
|
||||
streak_histogram = breaker.finalize()
|
||||
if streak_histogram:
|
||||
current_counters["block_streak_histogram"] = streak_histogram # type: ignore[assignment]
|
||||
runs_mod.mark_backfill_finished(
|
||||
db,
|
||||
run_id,
|
||||
|
|
|
|||
|
|
@ -62,6 +62,19 @@ def _make_snapshot(n: int) -> list[dict]:
|
|||
return [{"id": i + 1, "source_url": f"/items/{i + 1}"} for i in range(n)]
|
||||
|
||||
|
||||
def _fake_settings(**overrides: object) -> MagicMock:
|
||||
"""settings.* мокается атрибутами; bare MagicMock() без явных
|
||||
detail_backfill_block_ratio_window/_threshold отдаёт int(MagicMock())==1 --
|
||||
окно=1, порог=1.0 абортит на первом же блоке (#3184 review MINOR 4). Единая
|
||||
точка дефолтов вместо 23 копий тех же двух kwargs на каждом call site."""
|
||||
defaults: dict[str, object] = {
|
||||
"detail_backfill_block_ratio_window": 20,
|
||||
"detail_backfill_block_ratio_threshold": 0.7,
|
||||
}
|
||||
defaults.update(overrides)
|
||||
return MagicMock(**defaults)
|
||||
|
||||
|
||||
def _mock_db(snapshot: list[dict]) -> MagicMock:
|
||||
"""Fake Session: first execute() returns snapshot via .mappings().all()."""
|
||||
db = MagicMock()
|
||||
|
|
@ -99,7 +112,9 @@ async def test_backfill_empty_snapshot_marks_done() -> None:
|
|||
db = _mock_db([])
|
||||
runs = MagicMock()
|
||||
fake_session = AsyncMock()
|
||||
fake_settings = MagicMock(scraper_fetch_mode="cffi")
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=fake_session),
|
||||
|
|
@ -130,7 +145,9 @@ async def test_backfill_processes_snapshot_to_completion() -> None:
|
|||
mock_enrichment = MagicMock()
|
||||
mock_fetch = AsyncMock(return_value=mock_enrichment)
|
||||
mock_save = MagicMock(return_value=True)
|
||||
fake_settings = MagicMock(scraper_fetch_mode="cffi")
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
|
|
@ -178,7 +195,7 @@ async def test_backfill_build_warmed_session_receives_config() -> None:
|
|||
mock_enrichment = MagicMock()
|
||||
mock_fetch = AsyncMock(return_value=mock_enrichment)
|
||||
mock_save = MagicMock(return_value=True)
|
||||
fake_settings = MagicMock(
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_backfill_use_curl=True,
|
||||
)
|
||||
|
|
@ -225,7 +242,10 @@ async def test_backfill_reports_ban_kind_of_the_blocks_it_saw(
|
|||
runs = MagicMock()
|
||||
mock_scraper = MagicMock()
|
||||
mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True)
|
||||
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False)
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_backfill_use_curl=False,
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
|
|
@ -268,7 +288,10 @@ async def test_backfill_ban_kinds_census_keeps_multiplicities() -> None:
|
|||
)
|
||||
mock_scraper = MagicMock()
|
||||
mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True)
|
||||
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False)
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_backfill_use_curl=False,
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
|
|
@ -302,7 +325,10 @@ async def test_backfill_abort_log_has_no_ip_rate_limited_literal(caplog: Any) ->
|
|||
mock_fetch = AsyncMock(side_effect=AvitoBlockedError("firewall/soft-block"))
|
||||
mock_scraper = MagicMock()
|
||||
mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True)
|
||||
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False)
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_backfill_use_curl=False,
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
|
|
@ -344,7 +370,10 @@ async def test_backfill_blocked_abort_after_max_consecutive() -> None:
|
|||
mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True)
|
||||
# use_curl=False: legacy block→_rotate_ip путь (#1551 warm-batch rebuild только при
|
||||
# use_curl=True). MagicMock без явного флага сделал бы use_curl truthy.
|
||||
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False)
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_backfill_use_curl=False,
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
|
|
@ -383,7 +412,9 @@ async def test_backfill_sigterm_drain_breaks_and_marks_done_partial() -> None:
|
|||
mock_enrichment = MagicMock()
|
||||
mock_fetch = AsyncMock(return_value=mock_enrichment)
|
||||
mock_save = MagicMock(return_value=True)
|
||||
fake_settings = MagicMock(scraper_fetch_mode="cffi")
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
|
|
@ -417,7 +448,9 @@ async def test_backfill_budget_guard_stops_loop() -> None:
|
|||
db = _mock_db(snapshot)
|
||||
runs = MagicMock()
|
||||
mono_values = iter([0.0, 999.0, 999.0])
|
||||
fake_settings = MagicMock(scraper_fetch_mode="cffi")
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
|
|
@ -438,7 +471,9 @@ async def test_backfill_top_level_exception_marks_failed() -> None:
|
|||
db = MagicMock()
|
||||
db.execute.side_effect = RuntimeError("DB connection lost")
|
||||
runs = MagicMock()
|
||||
fake_settings = MagicMock(scraper_fetch_mode="cffi")
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
|
|
@ -469,7 +504,10 @@ async def test_backfill_rotate_ip_called_on_each_block() -> None:
|
|||
mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True)
|
||||
# use_curl=False: legacy block→_rotate_ip путь (#1551 warm-batch rebuild только при
|
||||
# use_curl=True). MagicMock без явного флага сделал бы use_curl truthy.
|
||||
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False)
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_backfill_use_curl=False,
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
|
|
@ -503,7 +541,9 @@ async def test_backfill_snapshot_filters_ekb_active_only() -> None:
|
|||
"""
|
||||
db = _mock_db([])
|
||||
runs = MagicMock()
|
||||
fake_settings = MagicMock(scraper_fetch_mode="cffi")
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
|
|
@ -649,7 +689,9 @@ async def test_backfill_fetch_exception_continues() -> None:
|
|||
runs = MagicMock()
|
||||
mock_enrichment = MagicMock()
|
||||
mock_fetch = AsyncMock(side_effect=[RuntimeError("parse error"), mock_enrichment])
|
||||
fake_settings = MagicMock(scraper_fetch_mode="cffi")
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
|
|
@ -700,7 +742,7 @@ async def test_backfill_fetch_timeout_skips_and_continues() -> None:
|
|||
return mock_enrichment
|
||||
|
||||
# Короткий timeout (50ms) — тест не ждёт реальные 90s; rotate-settle тоже 0.
|
||||
fake_settings = MagicMock(
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_fetch_timeout_s=0.05,
|
||||
avito_proxy_rotate_settle_s=0.0,
|
||||
|
|
@ -747,7 +789,10 @@ async def test_backfill_listing_gone_marks_inactive_no_breaker() -> None:
|
|||
mock_scraper.return_value._rotate_ip = AsyncMock(return_value=True)
|
||||
# use_curl=False: legacy block→_rotate_ip путь (#1551 warm-batch rebuild только при
|
||||
# use_curl=True). MagicMock без явного флага сделал бы use_curl truthy.
|
||||
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False)
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_backfill_use_curl=False,
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
|
|
@ -796,7 +841,7 @@ async def test_backfill_use_curl_flag_skips_browser_fetcher() -> None:
|
|||
mock_fetch = AsyncMock(return_value=mock_enrichment)
|
||||
mock_save = MagicMock(return_value=True)
|
||||
# Флаг use_curl=True, fetch_mode=browser (но флаг перекрывает)
|
||||
fake_settings = MagicMock(
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="browser",
|
||||
avito_detail_backfill_use_curl=True,
|
||||
)
|
||||
|
|
@ -835,7 +880,7 @@ async def test_backfill_use_curl_false_creates_browser_fetcher() -> None:
|
|||
mock_enrichment = MagicMock()
|
||||
mock_fetch = AsyncMock(return_value=mock_enrichment)
|
||||
mock_save = MagicMock(return_value=True)
|
||||
fake_settings = MagicMock(
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="browser",
|
||||
avito_detail_backfill_use_curl=False,
|
||||
)
|
||||
|
|
@ -900,7 +945,7 @@ async def test_backfill_use_curl_block_cooldown_research_no_rebuild() -> None:
|
|||
blocked_exc = AvitoBlockedError("rate-limited")
|
||||
mock_fetch = AsyncMock(side_effect=[blocked_exc, mock_enrichment])
|
||||
warmed_session = AsyncMock()
|
||||
fake_settings = MagicMock(
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_backfill_use_curl=True,
|
||||
)
|
||||
|
|
@ -963,7 +1008,10 @@ async def test_backfill_aborts_on_consecutive_failures_and_names_the_reason() ->
|
|||
"See https://curl.se/libcurl/c/libcurl-errors.html"
|
||||
)
|
||||
)
|
||||
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True)
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_backfill_use_curl=True,
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
|
|
@ -1003,7 +1051,10 @@ async def test_backfill_success_resets_failure_streak() -> None:
|
|||
boom = ValueError("avito detail HTTP 500 for https://www.avito.ru/x")
|
||||
# 2 отказа, успех, 2 отказа — при пороге 3 ни одна серия его не достигает.
|
||||
mock_fetch = AsyncMock(side_effect=[boom, boom, MagicMock(), boom, boom])
|
||||
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True)
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_backfill_use_curl=True,
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
|
|
@ -1022,3 +1073,238 @@ async def test_backfill_success_resets_failure_streak() -> None:
|
|||
assert result.attempted == 5
|
||||
assert result.enriched == 1
|
||||
assert result.failed == 4
|
||||
|
||||
|
||||
def _spaced_outcomes(total: int, block_positions: set[int]) -> list:
|
||||
"""1-indexed позиции block_positions -> AvitoBlockedError, остальные -- success MagicMock.
|
||||
|
||||
Ровный интервал между блоками (не пачками) -- моделирует прод-профиль 4786/5200
|
||||
(#3184): блоки есть, но локально в любом 20-окне их доля намного ниже порога.
|
||||
"""
|
||||
return [
|
||||
AvitoBlockedError("ip blocked") if pos in block_positions else MagicMock()
|
||||
for pos in range(1, total + 1)
|
||||
]
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_backfill_25pct_blocks_spread_does_not_abort() -> None:
|
||||
"""47 блоков из 190 попыток (25%, прогон 4786) -- НЕ обрывается (#3184).
|
||||
|
||||
Раньше это же соотношение (после единственной пачки >= max_consecutive_blocks)
|
||||
ловилось общим брейкером и уходило в 'banned', хотя прогон честно обогатил
|
||||
125/190 карточек. Блоки расставлены каждый 4-й (изолированно, не пачкой) --
|
||||
ни в одном 20-окне доля не приближается к порогу 0.7.
|
||||
"""
|
||||
total = 190
|
||||
block_positions = {pos for pos in range(1, total + 1) if pos % 4 == 0}
|
||||
assert len(block_positions) == 47
|
||||
snapshot = _make_snapshot(total)
|
||||
db = _mock_db(snapshot)
|
||||
runs = MagicMock()
|
||||
mock_fetch = AsyncMock(side_effect=_spaced_outcomes(total, block_positions))
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_backfill_use_curl=True,
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
patch(_SCRAPER),
|
||||
patch(_RUNS, runs),
|
||||
patch(_FETCH, mock_fetch),
|
||||
patch(_SAVE, return_value=True),
|
||||
patch(_SLEEP, new_callable=AsyncMock),
|
||||
):
|
||||
result = await run_avito_detail_backfill(
|
||||
db, run_id=101, params={"batch_size": total, "budget_sec": 3600}
|
||||
)
|
||||
|
||||
assert result.attempted == total
|
||||
assert result.blocked == 47
|
||||
assert result.enriched == total - 47
|
||||
runs.mark_backfill_finished.assert_called_once()
|
||||
assert runs.mark_backfill_finished.call_args.kwargs["aborted_by_blocks"] is False
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_backfill_32pct_blocks_spread_does_not_abort() -> None:
|
||||
"""20 блоков из 63 попыток (32%, прогон 5200) -- НЕ обрывается (#3184)."""
|
||||
total = 63
|
||||
block_positions = {pos for pos in range(3, 61, 3)}
|
||||
assert len(block_positions) == 20
|
||||
snapshot = _make_snapshot(total)
|
||||
db = _mock_db(snapshot)
|
||||
runs = MagicMock()
|
||||
mock_fetch = AsyncMock(side_effect=_spaced_outcomes(total, block_positions))
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_backfill_use_curl=True,
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
patch(_SCRAPER),
|
||||
patch(_RUNS, runs),
|
||||
patch(_FETCH, mock_fetch),
|
||||
patch(_SAVE, return_value=True),
|
||||
patch(_SLEEP, new_callable=AsyncMock),
|
||||
):
|
||||
result = await run_avito_detail_backfill(
|
||||
db, run_id=102, params={"batch_size": total, "budget_sec": 3600}
|
||||
)
|
||||
|
||||
assert result.attempted == total
|
||||
assert result.blocked == 20
|
||||
assert result.enriched == total - 20
|
||||
runs.mark_backfill_finished.assert_called_once()
|
||||
assert runs.mark_backfill_finished.call_args.kwargs["aborted_by_blocks"] is False
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_backfill_burst_mid_healthy_run_does_not_abort_and_streak_logged() -> None:
|
||||
"""6 блоков подряд посреди здорового прогона -- НЕ обрывает (#3184), а пачка
|
||||
|
||||
едет в counters["block_streak_histogram"]. Раньше max_consecutive_blocks=5
|
||||
абортил бы уже на 5-м блоке пачки независимо от контекста; серия сама по себе
|
||||
больше не критерий -- 20 успехов до пачки держат окно далеко от порога 0.7.
|
||||
"""
|
||||
total = 50
|
||||
block_positions = set(range(21, 27)) # 6 подряд, позиции 21..26
|
||||
assert len(block_positions) == 6
|
||||
snapshot = _make_snapshot(total)
|
||||
db = _mock_db(snapshot)
|
||||
runs = MagicMock()
|
||||
mock_fetch = AsyncMock(side_effect=_spaced_outcomes(total, block_positions))
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_backfill_use_curl=True,
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
patch(_SCRAPER),
|
||||
patch(_RUNS, runs),
|
||||
patch(_FETCH, mock_fetch),
|
||||
patch(_SAVE, return_value=True),
|
||||
patch(_SLEEP, new_callable=AsyncMock),
|
||||
):
|
||||
result = await run_avito_detail_backfill(
|
||||
db, run_id=103, params={"batch_size": total, "budget_sec": 3600}
|
||||
)
|
||||
|
||||
assert result.attempted == total
|
||||
assert result.blocked == 6
|
||||
runs.mark_backfill_finished.assert_called_once()
|
||||
assert runs.mark_backfill_finished.call_args.kwargs["aborted_by_blocks"] is False
|
||||
counters_arg = runs.mark_backfill_finished.call_args.args[2]
|
||||
assert counters_arg["block_streak_histogram"] == {"6": 1} # jsonb -- ключи строкой
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_backfill_ratio_criterion_aborts_after_window_fills_with_blocks() -> None:
|
||||
"""Успех на 1-й попытке, дальше сплошные блоки -- обрывает ПО ДОЛЕ в окне, не
|
||||
safety-net (#3184 review MAJOR 3: раньше ни один тест не проверял срабатывание
|
||||
ratio-критерия, только его отсутствие).
|
||||
|
||||
Окно (20) заполняется целиком ровно на 20-й попытке (1 успех + 19 блоков) --
|
||||
ratio=19/20=0.95 >= порога 0.7, абортит там же, до 21-й попытки дело не доходит.
|
||||
Safety-net здесь не может сработать: success на 1-й попытке снимает
|
||||
_pure_block_run сразу, до любого блока.
|
||||
"""
|
||||
total = 200
|
||||
block_positions = set(range(2, total + 1)) # всё, кроме 1-й попытки
|
||||
snapshot = _make_snapshot(total)
|
||||
db = _mock_db(snapshot)
|
||||
runs = MagicMock()
|
||||
mock_fetch = AsyncMock(side_effect=_spaced_outcomes(total, block_positions))
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_backfill_use_curl=True,
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
patch(_SCRAPER),
|
||||
patch(_RUNS, runs),
|
||||
patch(_FETCH, mock_fetch),
|
||||
patch(_SAVE, return_value=True),
|
||||
patch(_SLEEP, new_callable=AsyncMock),
|
||||
):
|
||||
result = await run_avito_detail_backfill(
|
||||
db, run_id=104, params={"batch_size": total, "budget_sec": 3600}
|
||||
)
|
||||
|
||||
assert result.attempted == 20
|
||||
assert result.blocked == 19
|
||||
assert result.enriched == 1
|
||||
runs.mark_backfill_finished.assert_called_once()
|
||||
assert runs.mark_backfill_finished.call_args.kwargs["aborted_by_blocks"] is True
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_backfill_ratio_boundary_14_of_20_aborts() -> None:
|
||||
"""14/20 = 0.7 -- ровно порог, should_abort сравнивает через >=, значит абортит
|
||||
(#3184 review MAJOR 3 -- граница снизу)."""
|
||||
total = 20
|
||||
block_positions = set(range(7, 21)) # позиции 7..20 включительно -- 14 штук
|
||||
assert len(block_positions) == 14
|
||||
snapshot = _make_snapshot(total)
|
||||
db = _mock_db(snapshot)
|
||||
runs = MagicMock()
|
||||
mock_fetch = AsyncMock(side_effect=_spaced_outcomes(total, block_positions))
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_backfill_use_curl=True,
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
patch(_SCRAPER),
|
||||
patch(_RUNS, runs),
|
||||
patch(_FETCH, mock_fetch),
|
||||
patch(_SAVE, return_value=True),
|
||||
patch(_SLEEP, new_callable=AsyncMock),
|
||||
):
|
||||
result = await run_avito_detail_backfill(
|
||||
db, run_id=105, params={"batch_size": total, "budget_sec": 3600}
|
||||
)
|
||||
|
||||
assert result.blocked == 14
|
||||
runs.mark_backfill_finished.assert_called_once()
|
||||
assert runs.mark_backfill_finished.call_args.kwargs["aborted_by_blocks"] is True
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_backfill_ratio_boundary_13_of_20_does_not_abort() -> None:
|
||||
"""13/20 = 0.65 -- ниже порога 0.7, НЕ абортит (#3184 review MAJOR 3 -- граница
|
||||
сверху, зеркало предыдущего теста)."""
|
||||
total = 20
|
||||
block_positions = set(range(8, 21)) # позиции 8..20 включительно -- 13 штук
|
||||
assert len(block_positions) == 13
|
||||
snapshot = _make_snapshot(total)
|
||||
db = _mock_db(snapshot)
|
||||
runs = MagicMock()
|
||||
mock_fetch = AsyncMock(side_effect=_spaced_outcomes(total, block_positions))
|
||||
fake_settings = _fake_settings(
|
||||
scraper_fetch_mode="cffi",
|
||||
avito_detail_backfill_use_curl=True,
|
||||
)
|
||||
with (
|
||||
patch(_SETTINGS, fake_settings),
|
||||
patch(_SESSION, return_value=AsyncMock()),
|
||||
patch(_SCRAPER),
|
||||
patch(_RUNS, runs),
|
||||
patch(_FETCH, mock_fetch),
|
||||
patch(_SAVE, return_value=True),
|
||||
patch(_SLEEP, new_callable=AsyncMock),
|
||||
):
|
||||
result = await run_avito_detail_backfill(
|
||||
db, run_id=106, params={"batch_size": total, "budget_sec": 3600}
|
||||
)
|
||||
|
||||
assert result.attempted == 20
|
||||
assert result.blocked == 13
|
||||
runs.mark_backfill_finished.assert_called_once()
|
||||
assert runs.mark_backfill_finished.call_args.kwargs["aborted_by_blocks"] is False
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue