feat(tradein/avito): ротация exit-IP по счётчику попыток в доборе карточек #3302

Merged
lekss361 merged 2 commits from feat/avito-rotate-by-attempt-counter into main 2026-08-31 12:09:04 +00:00
5 changed files with 437 additions and 0 deletions

View file

@ -1166,6 +1166,17 @@ class Settings(BaseSettings):
default=90.0, validation_alias="AVITO_DETAIL_FETCH_TIMEOUT_S"
)
# Ротация exit-IP по счётчику попыток на арендованном прокси (задача поверх #3298
# app.services.proxy_rotation.rotate_proxy). Замер вживую: билайновский порт держит
# ~15 карточек до полного обвала (13/39, последняя четверть 0%), настоящий
# мобильный оператор — ~30 (27/39, без обвала). Порог зависит от оператора порта,
# поэтому настройка, а не константа — подстройка под конкретный IP без релиза.
# Считаются ПОПЫТКИ (успех ИЛИ отказ), не только успехи: неудачная карточка тратит
# бюджет IP так же, как удачная. ENV: AVITO_DETAIL_BACKFILL_ROTATE_AFTER_ATTEMPTS.
avito_detail_backfill_rotate_after_attempts: int = Field(
default=15, ge=1, validation_alias="AVITO_DETAIL_BACKFILL_ROTATE_AFTER_ATTEMPTS"
)
# #3184: доля блоков в скользящем окне последних N попыток -- критерий обрыва
# avito_detail_backfill (app.services.backfill_block_breaker.BlockRatioBreaker),
# взамен голого "N блоков подряд". ТОЛЬКО avito -- изначальный план распространить

View file

@ -89,6 +89,7 @@ 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.proxy_rotation import rotate_proxy
from app.services.scraper_adapters import RealProxyProvider, RealScraperConfig
# #2397 Part D1 (#2330 закрыт): _AVITO_WARM_SEARCH_URL/build_warmed_session больше
@ -172,6 +173,75 @@ def _top_failure(census: Counter[str]) -> str | None:
return f"{reason} ({hits} из {sum(census.values())})"
# Провайдер mobileproxy.space почти всегда возвращает "rt" (секунды на переподключение
# канала после ротации) в RotationResult.reconnect_delay_s -- см. app.services.
# proxy_rotation docstring. Дефолт нужен ТОЛЬКО если провайдер его не прислал
# (best-effort парсинг "rt", неудача не является ошибкой ротации) -- замеренный
# вживую диапазон простоя канала 2-12с, 10с чуть выше нижней границы.
_DEFAULT_ROTATE_RECONNECT_DELAY_S = 10.0
async def _rotate_current_proxy(db: Session, run_id: int, browser_fetcher: BrowserFetcher) -> None:
"""Ротация exit-IP арендованного прокси по счётчику попыток + сброс browser-контекста.
Свежий адрес со старыми куками бесполезен личность браузера должна меняться
ВМЕСТЕ с адресом, иначе следующий запрос уходит с нового IP, но с cookie-следом
старого (request_context_reset потребляется ровно следующим fetch(), #3118).
Работает ТОЛЬКО когда есть browser_fetcher (browser-режим) это единственный
путь, где avito_detail_backfill реально держит lease с proxy_id (scrape_proxies.id,
browser_fetcher.lease_id). curl/backconnect-путь (use_curl=True) идёт через sticky
settings.scraper_proxy_url БЕЗ lease там ротировать нечего, вызывающий цикл не
зовёт эту функцию в том режиме вовсе.
Прод при этом в browser-режиме, а не в curl: у контейнера tradein-scraper (там же
живёт планировщик) проверено AVITO_DETAIL_BACKFILL_USE_CURL=false при
SCRAPER_FETCH_MODE=browser и USE_PROXY_POOL_BROWSER=true, так что ротация
активна. Значение true стоит только у tradein-backend, который добор не запускает.
Комментарий ниже по файлу (~строка 653) называет use_curl=True «прод-дефолтом»
это предсуществующее заблуждение, а не описание текущего прода.
Отказ провайдера (лимит исчерпан, нет rotate_url, сетевой сбой, неизвестный хост)
НЕ должен ронять прогон логируем и продолжаем на текущем адресе. Эта попытка НЕ
является блоком/отказом площадки и не должна попадать в BlockRatioBreaker/
counters.blocked/ban_kinds вызывающий код не передаёт её исход ни в один из них,
это обслуживание канала, а не результат fetch().
"""
proxy_id = browser_fetcher.lease_id
if proxy_id is None:
logger.info(
"avito_detail_backfill: run_id=%d rotate-by-attempts skipped -- no leased proxy",
run_id,
)
return
result = await rotate_proxy(db, proxy_id)
if not result.ok:
logger.warning(
"avito_detail_backfill: run_id=%d rotate-by-attempts FAILED proxy_id=%d: %s",
run_id,
proxy_id,
result.reason,
)
return
delay = (
result.reconnect_delay_s
if result.reconnect_delay_s is not None
else _DEFAULT_ROTATE_RECONNECT_DELAY_S
)
logger.info(
"avito_detail_backfill: run_id=%d rotate-by-attempts OK proxy_id=%d new_ip=%s -- "
"waiting %.1fs for channel reconnect",
run_id,
proxy_id,
result.new_ip,
delay,
)
await asyncio.sleep(delay)
browser_fetcher.request_context_reset()
@dataclass
class AvitoDetailBackfillResult:
"""Counters for one backfill run."""
@ -448,6 +518,11 @@ async def run_avito_detail_backfill(
abort_reason: str | None = None
do_sleep = False
items_since_warm = 0
# Счётчик попыток (успех ИЛИ отказ — оба тратят бюджет IP одинаково, см.
# settings.avito_detail_backfill_rotate_after_attempts) на ТЕКУЩЕМ арендованном
# прокси. Обнуляется на каждой ротации (успешной ИЛИ неуспешной — иначе
# исчерпанный дневной лимит провайдера дёргал бы rotate_proxy на каждой попытке).
attempts_since_rotation = 0
# #3251: сброс переиспользуемого browser-context'а (reuse_context=True выше)
# разрешён РОВНО один раз за прогон — зеркалит domclick_detail_backfill (#3212).
# Сброс на КАЖДЫЙ блок сам себя поддерживает: пройденный QRATOR-PoW живёт в
@ -502,6 +577,7 @@ async def run_avito_detail_backfill(
source_url: str = row["source_url"]
counters.attempted += 1
attempts_since_rotation += 1
# Normalise URL -> path (mirrors run_avito_city_sweep detail-phase, kit)
item_url = urlparse(source_url).path if source_url.startswith("http") else source_url
@ -769,6 +845,18 @@ async def run_avito_detail_backfill(
except Exception:
pass
# Ротация exit-IP по счётчику попыток на арендованном прокси (задача поверх
# #3298). browser_fetcher is not None -- см. _rotate_current_proxy docstring
# за тем, почему только browser-режим держит lease с proxy_id. Обнуляем
# счётчик ДО вызова (не после) -- неудача ротации не должна дёргать
# rotate_proxy на КАЖДОЙ следующей попытке до конца прогона.
if (
browser_fetcher is not None
and attempts_since_rotation >= settings.avito_detail_backfill_rotate_after_attempts
):
attempts_since_rotation = 0
await _rotate_current_proxy(db, run_id, browser_fetcher)
if consecutive_failures >= max_consecutive_failures:
logger.error(
"avito_detail_backfill: run_id=%d ABORT -- %d отказов подряд без "

View file

@ -19,7 +19,9 @@ from scraper_kit.avito_exceptions import ( # noqa: E402
)
from app.core import shutdown as _sd # noqa: E402
from app.services.proxy_rotation import RotationResult # noqa: E402
from app.tasks.avito_detail_backfill import ( # noqa: E402
_DEFAULT_ROTATE_RECONNECT_DELAY_S,
_OBLAST_AVITO_URL_PATTERNS,
AvitoDetailBackfillResult,
run_avito_detail_backfill,
@ -70,6 +72,14 @@ def _fake_settings(**overrides: object) -> MagicMock:
defaults: dict[str, object] = {
"detail_backfill_block_ratio_window": 20,
"detail_backfill_block_ratio_threshold": 0.7,
# Реальный дефолт (config.py) -- без него bare MagicMock() возвращает
# child-MagicMock на сравнение `attempts_since_rotation >= settings.avito_
# detail_backfill_rotate_after_attempts` в browser-режиме и падает TypeError
# (int >= MagicMock не поддерживается). В curl-режиме (browser_fetcher=None)
# сравнение вообще не вычисляется (short-circuit `and`), но browser-тесты
# ниже (test_backfill_use_curl_false_creates_browser_fetcher и rotate-тесты)
# его достигают.
"avito_detail_backfill_rotate_after_attempts": 15,
}
defaults.update(overrides)
return MagicMock(**defaults)
@ -100,6 +110,7 @@ _RESEARCH = "app.tasks.avito_detail_backfill.research_in_session"
# тесты про block/ban/rotate-логику, не про подбор прокси (см.
# tests/services/test_proxy_egress.py), поэтому мокаем сам резолвер.
_RESOLVE_PROXY_URL = "app.tasks.avito_detail_backfill.resolve_proxy_url"
_ROTATE_PROXY = "app.tasks.avito_detail_backfill.rotate_proxy"
# ---------------------------------------------------------------------------
# Tests
@ -1431,3 +1442,313 @@ async def test_backfill_no_abort_leaves_no_abort_reason_in_counters() -> None:
counters = runs.mark_backfill_finished.call_args.args[2]
assert "abort_reason" not in counters
# ---------------------------------------------------------------------------
# Rotate-by-attempt-counter (задача поверх #3298)
# ---------------------------------------------------------------------------
def _browser_mock_instance(lease_id: int | None) -> AsyncMock:
"""Мокнутый BrowserFetcher instance с lease_id и sync request_context_reset
(реальный метод -- НЕ async, mock должен звать его так же)."""
inst = AsyncMock()
inst.__aenter__ = AsyncMock(return_value=inst)
inst.__aexit__ = AsyncMock(return_value=False)
inst.lease_id = lease_id
inst.request_context_reset = MagicMock()
return inst
@pytest.mark.asyncio
async def test_rotate_triggers_at_threshold() -> None:
"""attempts_since_rotation достигает порога -> rotate_proxy вызван РОВНО один раз
с proxy_id текущего lease (browser_fetcher.lease_id)."""
snapshot = _make_snapshot(3)
db = _mock_db(snapshot)
runs = MagicMock()
mock_fetch = AsyncMock(return_value=MagicMock())
mock_bf_instance = _browser_mock_instance(lease_id=777)
fake_settings = _fake_settings(
scraper_fetch_mode="browser",
avito_detail_backfill_use_curl=False,
avito_detail_backfill_rotate_after_attempts=3,
)
mock_rotate = AsyncMock(
return_value=RotationResult(ok=True, reason=None, new_ip="1.2.3.4", reconnect_delay_s=4.0)
)
with (
patch(_SETTINGS, fake_settings),
patch(_SESSION),
patch(_SCRAPER),
patch(_BROWSER_FETCHER, return_value=mock_bf_instance),
patch(_RUNS, runs),
patch(_FETCH, mock_fetch),
patch(_SAVE, return_value=True),
patch(_SLEEP, new_callable=AsyncMock),
patch(_ROTATE_PROXY, mock_rotate),
):
await run_avito_detail_backfill(
db, run_id=200, params={"batch_size": 3, "budget_sec": 3600}
)
mock_rotate.assert_called_once()
args, _ = mock_rotate.call_args
assert args[1] == 777 # proxy_id из lease_id
@pytest.mark.asyncio
async def test_rotate_resets_counter_after_trigger() -> None:
"""Порог=2, снапшот=4 -> ротация срабатывает ДВАЖДЫ (счётчик обнуляется после
каждого срабатывания, не только один раз за прогон)."""
snapshot = _make_snapshot(4)
db = _mock_db(snapshot)
runs = MagicMock()
mock_fetch = AsyncMock(return_value=MagicMock())
mock_bf_instance = _browser_mock_instance(lease_id=5)
fake_settings = _fake_settings(
scraper_fetch_mode="browser",
avito_detail_backfill_use_curl=False,
avito_detail_backfill_rotate_after_attempts=2,
)
mock_rotate = AsyncMock(
return_value=RotationResult(ok=True, reason=None, new_ip="9.9.9.9", reconnect_delay_s=1.0)
)
with (
patch(_SETTINGS, fake_settings),
patch(_SESSION),
patch(_SCRAPER),
patch(_BROWSER_FETCHER, return_value=mock_bf_instance),
patch(_RUNS, runs),
patch(_FETCH, mock_fetch),
patch(_SAVE, return_value=True),
patch(_SLEEP, new_callable=AsyncMock),
patch(_ROTATE_PROXY, mock_rotate),
):
await run_avito_detail_backfill(
db, run_id=201, params={"batch_size": 4, "budget_sec": 3600}
)
assert mock_rotate.call_count == 2
@pytest.mark.asyncio
async def test_rotate_waits_reconnect_delay_from_provider() -> None:
"""reconnect_delay_s из RotationResult -- ждём именно его, не дефолт."""
snapshot = _make_snapshot(1)
db = _mock_db(snapshot)
runs = MagicMock()
mock_fetch = AsyncMock(return_value=MagicMock())
mock_bf_instance = _browser_mock_instance(lease_id=1)
fake_settings = _fake_settings(
scraper_fetch_mode="browser",
avito_detail_backfill_use_curl=False,
avito_detail_backfill_rotate_after_attempts=1,
)
mock_rotate = AsyncMock(
return_value=RotationResult(ok=True, reason=None, new_ip="1.1.1.1", reconnect_delay_s=6.5)
)
mock_sleep = AsyncMock()
with (
patch(_SETTINGS, fake_settings),
patch(_SESSION),
patch(_SCRAPER),
patch(_BROWSER_FETCHER, return_value=mock_bf_instance),
patch(_RUNS, runs),
patch(_FETCH, mock_fetch),
patch(_SAVE, return_value=True),
patch(_SLEEP, mock_sleep),
patch(_ROTATE_PROXY, mock_rotate),
):
await run_avito_detail_backfill(
db, run_id=202, params={"batch_size": 1, "budget_sec": 3600}
)
sleep_values = [call.args[0] for call in mock_sleep.call_args_list]
assert 6.5 in sleep_values
@pytest.mark.asyncio
async def test_rotate_uses_default_delay_when_provider_silent() -> None:
"""reconnect_delay_s=None (провайдер не прислал rt) -> используем дефолт
_DEFAULT_ROTATE_RECONNECT_DELAY_S, не падаем и не ждём 0с."""
snapshot = _make_snapshot(1)
db = _mock_db(snapshot)
runs = MagicMock()
mock_fetch = AsyncMock(return_value=MagicMock())
mock_bf_instance = _browser_mock_instance(lease_id=1)
fake_settings = _fake_settings(
scraper_fetch_mode="browser",
avito_detail_backfill_use_curl=False,
avito_detail_backfill_rotate_after_attempts=1,
)
mock_rotate = AsyncMock(
return_value=RotationResult(ok=True, reason=None, new_ip="2.2.2.2", reconnect_delay_s=None)
)
mock_sleep = AsyncMock()
with (
patch(_SETTINGS, fake_settings),
patch(_SESSION),
patch(_SCRAPER),
patch(_BROWSER_FETCHER, return_value=mock_bf_instance),
patch(_RUNS, runs),
patch(_FETCH, mock_fetch),
patch(_SAVE, return_value=True),
patch(_SLEEP, mock_sleep),
patch(_ROTATE_PROXY, mock_rotate),
):
await run_avito_detail_backfill(
db, run_id=203, params={"batch_size": 1, "budget_sec": 3600}
)
sleep_values = [call.args[0] for call in mock_sleep.call_args_list]
assert _DEFAULT_ROTATE_RECONNECT_DELAY_S in sleep_values
@pytest.mark.asyncio
async def test_rotate_success_calls_context_reset() -> None:
"""Успешная ротация -> request_context_reset() вызван (новый адрес + свежий
browser-контекст, старые куки выброшены)."""
snapshot = _make_snapshot(1)
db = _mock_db(snapshot)
runs = MagicMock()
mock_fetch = AsyncMock(return_value=MagicMock())
mock_bf_instance = _browser_mock_instance(lease_id=42)
fake_settings = _fake_settings(
scraper_fetch_mode="browser",
avito_detail_backfill_use_curl=False,
avito_detail_backfill_rotate_after_attempts=1,
)
mock_rotate = AsyncMock(
return_value=RotationResult(ok=True, reason=None, new_ip="3.3.3.3", reconnect_delay_s=1.0)
)
with (
patch(_SETTINGS, fake_settings),
patch(_SESSION),
patch(_SCRAPER),
patch(_BROWSER_FETCHER, return_value=mock_bf_instance),
patch(_RUNS, runs),
patch(_FETCH, mock_fetch),
patch(_SAVE, return_value=True),
patch(_SLEEP, new_callable=AsyncMock),
patch(_ROTATE_PROXY, mock_rotate),
):
await run_avito_detail_backfill(
db, run_id=204, params={"batch_size": 1, "budget_sec": 3600}
)
mock_bf_instance.request_context_reset.assert_called_once()
@pytest.mark.asyncio
async def test_rotate_failure_does_not_abort_run_and_skips_context_reset() -> None:
"""rotate_proxy(ok=False) (лимит исчерпан / нет rotate_url / провайдер отказал)
-- прогон НЕ падает, продолжает на текущем адресе, request_context_reset НЕ
вызывается (менять контекст без нового IP бессмысленно)."""
snapshot = _make_snapshot(2)
db = _mock_db(snapshot)
runs = MagicMock()
mock_fetch = AsyncMock(return_value=MagicMock())
mock_save = MagicMock(return_value=True)
mock_bf_instance = _browser_mock_instance(lease_id=9)
fake_settings = _fake_settings(
scraper_fetch_mode="browser",
avito_detail_backfill_use_curl=False,
avito_detail_backfill_rotate_after_attempts=1,
)
mock_rotate = AsyncMock(
return_value=RotationResult(ok=False, reason="daily rotation limit reached (200/day)")
)
with (
patch(_SETTINGS, fake_settings),
patch(_SESSION),
patch(_SCRAPER),
patch(_BROWSER_FETCHER, return_value=mock_bf_instance),
patch(_RUNS, runs),
patch(_FETCH, mock_fetch),
patch(_SAVE, mock_save),
patch(_SLEEP, new_callable=AsyncMock),
patch(_ROTATE_PROXY, mock_rotate),
):
result = await run_avito_detail_backfill(
db, run_id=205, params={"batch_size": 2, "budget_sec": 3600}
)
# прогон не упал и довёл оба листинга до конца
assert result.enriched == 2
assert mock_rotate.call_count == 2 # счётчик обнулился и после неудачи -> 2 попытки
mock_bf_instance.request_context_reset.assert_not_called()
runs.mark_backfill_finished.assert_called_once()
# rotate НЕ считается блоком площадки
assert result.blocked == 0
@pytest.mark.asyncio
async def test_rotate_not_counted_as_platform_block() -> None:
"""Ротация не трогает breaker/blocked/ban_kinds -- это обслуживание канала, не
исход fetch(). aborted_by_blocks остаётся False, blocked=0 даже когда rotate
срабатывает на каждой попытке."""
snapshot = _make_snapshot(5)
db = _mock_db(snapshot)
runs = MagicMock()
mock_fetch = AsyncMock(return_value=MagicMock())
mock_bf_instance = _browser_mock_instance(lease_id=3)
fake_settings = _fake_settings(
scraper_fetch_mode="browser",
avito_detail_backfill_use_curl=False,
avito_detail_backfill_rotate_after_attempts=1,
)
mock_rotate = AsyncMock(
return_value=RotationResult(ok=True, reason=None, new_ip="4.4.4.4", reconnect_delay_s=0.1)
)
with (
patch(_SETTINGS, fake_settings),
patch(_SESSION),
patch(_SCRAPER),
patch(_BROWSER_FETCHER, return_value=mock_bf_instance),
patch(_RUNS, runs),
patch(_FETCH, mock_fetch),
patch(_SAVE, return_value=True),
patch(_SLEEP, new_callable=AsyncMock),
patch(_ROTATE_PROXY, mock_rotate),
):
result = await run_avito_detail_backfill(
db, run_id=206, params={"batch_size": 5, "budget_sec": 3600}
)
assert mock_rotate.call_count == 5
assert result.blocked == 0
finished_kwargs = runs.mark_backfill_finished.call_args.kwargs
assert finished_kwargs.get("aborted_by_blocks") is False
assert not finished_kwargs.get("ban_kinds")
@pytest.mark.asyncio
async def test_rotate_skipped_in_curl_mode_no_lease() -> None:
"""use_curl=True (прод-дефолт): browser_fetcher остаётся None -- rotate_proxy НЕ
вызывается вообще (curl/backconnect путь не держит lease с proxy_id)."""
snapshot = _make_snapshot(20)
db = _mock_db(snapshot)
runs = MagicMock()
mock_fetch = AsyncMock(return_value=MagicMock())
fake_settings = _fake_settings(
scraper_fetch_mode="cffi",
avito_detail_backfill_use_curl=True,
avito_detail_backfill_rotate_after_attempts=1,
)
mock_rotate = AsyncMock()
with (
patch(_SETTINGS, fake_settings),
patch(_SESSION),
patch(_SCRAPER),
patch(_RUNS, runs),
patch(_FETCH, mock_fetch),
patch(_SAVE, return_value=True),
patch(_SLEEP, new_callable=AsyncMock),
patch(_ROTATE_PROXY, mock_rotate),
):
await run_avito_detail_backfill(
db, run_id=207, params={"batch_size": 20, "budget_sec": 3600}
)
mock_rotate.assert_not_called()

View file

@ -67,6 +67,12 @@ def _fake_settings(**overrides: object) -> MagicMock:
"detail_backfill_block_ratio_window": 20,
"detail_backfill_block_ratio_threshold": 0.7,
"browser_http_endpoint": "http://browser:9000",
# Реальный дефолт (config.py). Без него bare MagicMock отдаёт
# child-MagicMock на сравнение `attempts_since_rotation >= settings.avito_
# detail_backfill_rotate_after_attempts` в browser-режиме, и тест падает
# TypeError: '>=' not supported between 'int' and 'MagicMock'. Та же
# ловушка, что уже описана здесь для detail_backfill_block_ratio_window.
"avito_detail_backfill_rotate_after_attempts": 15,
}
defaults.update(overrides)
return MagicMock(**defaults)

View file

@ -741,6 +741,17 @@ class BrowserFetcher:
return None, None
return self._lease.url, self._lease.kind
@property
def lease_id(self) -> int | None:
"""id ТЕКУЩЕГО session-lease (`scrape_proxies.id`) — None, если lease нет
(env-fallback путь, use_pool=False, или proxy_provider не подключён).
Минимальный read-only доступ для вызывающего кода, которому нужен proxy_id
арендованного прокси (например, для app.services.proxy_rotation.rotate_proxy)
БЕЗ прямого доступа к приватному `_lease` и без изменения lease-механики.
"""
return self._lease.id if self._lease is not None else None
def _report_fetch_result(self, ok: bool) -> None:
"""Учесть исход ОДНОГО /fetch в здоровье текущего session-lease.