fix(tradein/scraper): пустой пул прокси перестаёт стирать чекпоинт полного обхода #2834

Merged
bot-backend merged 1 commit from fix/2687-full-load-infra into main 2026-08-12 12:18:03 +00:00
2 changed files with 218 additions and 1 deletions

View file

@ -0,0 +1,166 @@
"""#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,
proxy_rotate_attempts=1,
proxy_rotate_attempt_timeout_s=1.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"]

View file

@ -54,6 +54,7 @@ from scraper_kit.avito_exceptions import (
from scraper_kit.base import ScrapedLot, save_listings
from scraper_kit.browser_fetcher import BrowserFetcher
from scraper_kit.orchestration import runs
# Константы диагноза берём напрямую, а не через `runs.` — helper ниже обязан
# работать и когда тесты подменяют весь модуль runs двойником (#2686).
from scraper_kit.orchestration.runs import BAN_KIND_INFRA, BAN_KIND_PLATFORM, BAN_KIND_UNKNOWN
@ -76,6 +77,7 @@ from scraper_kit.providers.yandex.serp import (
ROOM_PATH,
YandexRealtyScraper,
)
from scraper_kit.proxy_errors import NoProxyAvailableError
if TYPE_CHECKING:
from sqlalchemy.orm import Session
@ -208,7 +210,11 @@ def ban_kind_of_exception(exc: BaseException) -> str:
Один узел на все avito-сайты mark_banned: разводить исход в каждом из трёх
было бы тремя копиями одного условия.
"""
if isinstance(exc, AvitoSidecarUnavailableError):
if isinstance(exc, AvitoSidecarUnavailableError | NoProxyAvailableError):
# NoProxyAvailableError (#2687): пул прокси пуст в проде. Причина здесь не
# «не установлена» (тогда был бы 'unknown', #2764), а названа самим типом —
# см. proxy_errors.py: «это НАША инфраструктура (нет живого прокси), не
# внешний блок». Тот же диагноз, что и у отказа сайдкара.
return BAN_KIND_INFRA
if isinstance(exc, AvitoBlockedError | AvitoRateLimitedError):
return BAN_KIND_PLATFORM
@ -3217,6 +3223,20 @@ async def run_cian_full_load(
)
return counters
except NoProxyAvailableError as exc:
# #2687, см. ту же ветку в run_avito_full_load: наш отказ (пул пуст) обязан
# сохранять done_buckets, а generic-RuntimeError → mark_failed его теряет.
logger.error("cian-full-load run_id=%d: no proxy available — %s", run_id, exc)
counters.errors_count += 1
runs.mark_banned(
db,
run_id,
f"cian full load aborted: {exc}",
{**counters.to_dict(), "done_buckets": sorted(done)},
ban_kind=ban_kind_of_exception(exc),
)
return counters
except RuntimeError as exc:
# on_bucket кидает RuntimeError("cancelled") при cooperative cancel,
# RuntimeError("shutdown") при кооперативном SIGTERM-drain (#1182 Phase 3a).
@ -3445,6 +3465,20 @@ async def run_yandex_full_load(
)
return counters
except NoProxyAvailableError as exc:
# #2687, см. ту же ветку в run_avito_full_load: наш отказ (пул пуст) обязан
# сохранять done_buckets, а generic-RuntimeError → mark_failed его теряет.
logger.error("yandex-full-load run_id=%d: no proxy available — %s", run_id, exc)
counters.errors_count += 1
runs.mark_banned(
db,
run_id,
f"yandex full load aborted: {exc}",
{**counters.to_dict(), "done_buckets": sorted(done)},
ban_kind=ban_kind_of_exception(exc),
)
return counters
except RuntimeError as exc:
# on_bucket кидает RuntimeError("cancelled") при cooperative cancel,
# RuntimeError("shutdown") при кооперативном SIGTERM-drain (#1182 Phase 3a).
@ -3662,6 +3696,23 @@ async def run_avito_full_load(
)
return counters
except NoProxyAvailableError as exc:
# #2687: пул опустел mid-run. Ветка стоит ДО generic-RuntimeError намеренно —
# NoProxyAvailableError его подкласс, и без неё отказ уходил в mark_failed,
# который (в отличие от mark_banned) НЕ пишет done_buckets. То есть чекпоинт
# терялся ровно на НАШЕМ отказе — том исходе, для которого #2686 требовал его
# сохранять наравне с блокировкой площадкой.
logger.error("avito-full-load run_id=%d: no proxy available — %s", run_id, exc)
counters.errors_count += 1
runs.mark_banned(
db,
run_id,
f"avito full load aborted: {exc}",
{**counters.to_dict(), "done_buckets": sorted(done)},
ban_kind=ban_kind_of_exception(exc),
)
return counters
except RuntimeError as exc:
# on_bucket кидает RuntimeError("cancelled") при cooperative cancel,
# RuntimeError("shutdown") при кооперативном SIGTERM-drain (#1182 Phase 3a).