All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 11s
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 5m8s
proxy_rotate_attempts / proxy_rotate_attempt_timeout_s тюнили ретраи changeip-GET. Сам changeip снят в #2616 шаг 2 (аккаунт mobileproxy закрыт, ссылки нет), и с тех пор ручки живут пламбингом: Settings -> property адаптера -> поле протокола ScraperConfig -> и всё. Ни одного потребителя, только четыре теста, которые заполняют их при сборке конфига. Комментарий над ними в contracts.py оправдывал их сохранение так: «оставлены как budget-верхняя-граница для app.tasks.avito_detail_backfill wait_for» — но wait_for там берёт СОСЕДНЕЕ поле, avito_proxy_rotate_settle_s (avito_detail_backfill.py:249). То есть комментарий приписывал этим двум полям работу третьего и тем самым прикрывал их мёртвость. Соседние ручки проверены и ОСТАВЛЕНЫ, они действительно читаются: * avito/cian/yandex_proxy_max_rotations — pipeline._max_rotations; * avito_proxy_rotate_settle_s — asyncio.wait_for в avito_detail_backfill. Заодно сжат комментарий в config.py: перечисление истории changeip заменено на то, что нужно знать сейчас — кто читает оставшиеся две ручки и где живая ротация (ASOCKS_API_TOKEN / proxy_rotation, #2611). ruff clean, 4974 passed / 37 skipped.
164 lines
7.7 KiB
Python
164 lines
7.7 KiB
Python
"""#2687: опустевший пул прокси mid-run стирал чекпоинт полного обхода.
|
||
|
||
Найдено при разборе шести `ban_kind='infra'` у `avito_full_load` (29.07-03.08). Сами
|
||
шесть объяснены и починены раньше (#2637 подключил браузерный путь Авито к пулу,
|
||
#2634/#2640 добили запасной путь), но рядом с ними на main живёт соседний дефект того
|
||
же класса, что и #2686: наш собственный отказ финализируется веткой, которая теряет
|
||
чекпоинт.
|
||
|
||
Механизм. `NoProxyAvailableError` (proxy_errors.py) — наследник `RuntimeError`, и его
|
||
собственный докстринг говорит прямо: «это НАША инфраструктура (нет живого прокси), не
|
||
внешний блок». В `run_*_full_load` он попадал в общую ветку `except RuntimeError` →
|
||
`mark_failed`, а `mark_failed`, в отличие от `mark_banned`, НЕ пишет `done_buckets`
|
||
(pipeline.py, ветка блока пишет его явно). Плюс `ban_kind_of_exception` возвращал для
|
||
него `unknown` — диагноз «наша инфраструктура» терялся дважды: и как метка, и как
|
||
сохранённый прогресс.
|
||
|
||
Достижимость на проде (замер 2026-08-12, read-only):
|
||
- под `avito` доступны три узла (id 9/10/11; id 1 забанен парой до 13.08);
|
||
- `BrowserFetcher` меняет lease после `_LEASE_ROTATE_AFTER_FAILS=3` подряд неудачных
|
||
`/fetch`, а `proxy_pool.acquire` не выдаёт узел с
|
||
`consecutive_fails >= MAX_CONSECUTIVE_FAILS=3` → девять неудачных POST'ов
|
||
опустошают пул под источник;
|
||
- лестница ОДНОГО полного обхода при отказывающем сайдкаре делает до 24 POST'ов
|
||
(4 бакета `_AVITO_SWEEP_MAX_CONSECUTIVE_BLOCKED` × 3 попытки
|
||
`_AVITO_SIDECAR_TRANSIENT_RETRIES`+1 × 2 внутренних httpx-ретрая `fetch`),
|
||
то есть заведомо проходит через это состояние.
|
||
Цена потери чекпоинта — прогон 3547 (09.08): 35 бакетов, 5496 объявлений за 2ч58м.
|
||
|
||
Фальсификация (красный прогон на коде до правки):
|
||
- `ban_kind_of_exception` → `'unknown'` вместо `'infra'`;
|
||
- `run_*_full_load` → `mark_failed` вместо `mark_banned`, исключение улетает наружу
|
||
(тест падает на непойманном `NoProxyAvailableError`), `done_buckets` нигде нет.
|
||
"""
|
||
|
||
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 (
|
||
ban_kind_of_exception,
|
||
run_avito_full_load,
|
||
run_cian_full_load,
|
||
run_yandex_full_load,
|
||
)
|
||
from scraper_kit.proxy_errors import NoProxyAvailableError
|
||
|
||
PFX = "scraper_kit.orchestration.pipeline"
|
||
|
||
_FULL_LOAD = {
|
||
"avito": (run_avito_full_load, f"{PFX}.AvitoScraper"),
|
||
"cian": (run_cian_full_load, f"{PFX}.CianScraper"),
|
||
"yandex": (run_yandex_full_load, f"{PFX}.YandexRealtyScraper"),
|
||
}
|
||
|
||
|
||
class _Recorder:
|
||
"""Двойник scrape_runs: интересны имя финализатора, ban_kind и counters."""
|
||
|
||
def __init__(self) -> None:
|
||
self.calls: list[tuple[str, str, dict[str, Any]]] = []
|
||
|
||
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:
|
||
pass
|
||
|
||
def mark_done(self, db: Any, run_id: int, counters: dict[str, Any]) -> None:
|
||
self.calls.append(("mark_done", "", dict(counters)))
|
||
|
||
def mark_failed(self, db: Any, run_id: int, error: str, counters: dict[str, Any]) -> None:
|
||
self.calls.append(("mark_failed", "", dict(counters)))
|
||
|
||
def mark_banned(
|
||
self,
|
||
db: Any,
|
||
run_id: int,
|
||
error: str,
|
||
counters: dict[str, Any],
|
||
*,
|
||
ban_kind: str = kit_runs.BAN_KIND_UNKNOWN,
|
||
) -> None:
|
||
self.calls.append(("mark_banned", ban_kind, dict(counters)))
|
||
|
||
|
||
def _config() -> SimpleNamespace:
|
||
return SimpleNamespace(
|
||
scraper_fetch_mode="curl_cffi",
|
||
browser_http_endpoint="http://browser.test/fetch",
|
||
scraper_proxy_url=None,
|
||
avito_proxy_max_rotations=0,
|
||
avito_serp_ok_not_banned=True,
|
||
avito_proxy_rotate_settle_s=0.0,
|
||
cian_proxy_max_rotations=0,
|
||
yandex_proxy_max_rotations=0,
|
||
cian_full_load_per_fetch_timeout_s=0.0,
|
||
scraper_skip_seen_today=False,
|
||
)
|
||
|
||
|
||
def _scraper_that_saves_one_bucket_then_runs_out_of_proxies() -> MagicMock:
|
||
"""Двойник скрапера: один бакет доехал до on_bucket, затем пул опустел."""
|
||
|
||
async def _fetch(*_a: Any, on_bucket: Any = None, **_k: Any) -> None:
|
||
on_bucket("2к:0-5m", [MagicMock(source_id="a1")])
|
||
raise NoProxyAvailableError("avito")
|
||
|
||
m = MagicMock()
|
||
m.__aenter__ = AsyncMock(return_value=m)
|
||
m.__aexit__ = AsyncMock(return_value=None)
|
||
m.fetch_all_secondary = _fetch
|
||
m.state_extraction_attempts = 0
|
||
m.state_extraction_failures = 0
|
||
m.gate_fetch_attempts = 0
|
||
m.gate_fetch_failures = 0
|
||
m._browser = None
|
||
return m
|
||
|
||
|
||
def test_ban_kind_of_no_proxy_is_infra() -> None:
|
||
"""Пустой пул — наша инфраструктура по определению самого исключения.
|
||
|
||
Раньше классификатор отдавал 'unknown' (дефолт #2764): для НЕизвестной причины
|
||
это верно, но здесь причина известна и названа в типе.
|
||
"""
|
||
assert ban_kind_of_exception(NoProxyAvailableError("avito")) == kit_runs.BAN_KIND_INFRA
|
||
|
||
|
||
@pytest.mark.parametrize("source", ["avito", "cian", "yandex"])
|
||
async def test_full_load_pool_exhaustion_keeps_checkpoint(source: str) -> None:
|
||
"""Пул опустел mid-run → 'banned'/'infra' и done_buckets целы, а не mark_failed.
|
||
|
||
Требование #2686 («чекпоинт сохраняется и при нашем сбое, и при блокировке
|
||
площадкой») до этой правки держалось только для отказа сайдкара; отказ пула шёл
|
||
мимо него и обнулял прогресс следующего прогона.
|
||
"""
|
||
fn, scraper_target = _FULL_LOAD[source]
|
||
recorder = _Recorder()
|
||
extra: dict[str, Any] = {}
|
||
if source == "yandex":
|
||
enrichment = MagicMock()
|
||
enrichment.record_yandex_price_history = MagicMock(return_value=0)
|
||
extra["enrichment"] = enrichment
|
||
|
||
scraper = _scraper_that_saves_one_bucket_then_runs_out_of_proxies()
|
||
with (
|
||
patch(scraper_target, return_value=scraper),
|
||
patch(f"{PFX}.save_listings", MagicMock(return_value=(1, 0))),
|
||
patch(f"{PFX}.runs", recorder),
|
||
):
|
||
await fn(MagicMock(), run_id=1, config=_config(), matcher=MagicMock(), **extra)
|
||
|
||
assert [c[0] for c in recorder.calls] == ["mark_banned"]
|
||
_method, ban_kind, counters = recorder.calls[0]
|
||
assert ban_kind == kit_runs.BAN_KIND_INFRA
|
||
assert counters["done_buckets"] == ["2к:0-5m"]
|