fix(tradein/scraper): пустой пул прокси перестаёт стирать чекпоинт полного обхода (#2834)
All checks were successful
Deploy Trade-In / changes (push) Successful in 10s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m6s
Deploy Trade-In / build-backend (push) Successful in 1m36s
Deploy Trade-In / deploy (push) Successful in 1m23s
All checks were successful
Deploy Trade-In / changes (push) Successful in 10s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m6s
Deploy Trade-In / build-backend (push) Successful in 1m36s
Deploy Trade-In / deploy (push) Successful in 1m23s
This commit is contained in:
parent
e17687aed7
commit
4d31a0ee82
2 changed files with 218 additions and 1 deletions
|
|
@ -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"]
|
||||||
|
|
@ -54,6 +54,7 @@ from scraper_kit.avito_exceptions import (
|
||||||
from scraper_kit.base import ScrapedLot, save_listings
|
from scraper_kit.base import ScrapedLot, save_listings
|
||||||
from scraper_kit.browser_fetcher import BrowserFetcher
|
from scraper_kit.browser_fetcher import BrowserFetcher
|
||||||
from scraper_kit.orchestration import runs
|
from scraper_kit.orchestration import runs
|
||||||
|
|
||||||
# Константы диагноза берём напрямую, а не через `runs.` — helper ниже обязан
|
# Константы диагноза берём напрямую, а не через `runs.` — helper ниже обязан
|
||||||
# работать и когда тесты подменяют весь модуль runs двойником (#2686).
|
# работать и когда тесты подменяют весь модуль runs двойником (#2686).
|
||||||
from scraper_kit.orchestration.runs import BAN_KIND_INFRA, BAN_KIND_PLATFORM, BAN_KIND_UNKNOWN
|
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,
|
ROOM_PATH,
|
||||||
YandexRealtyScraper,
|
YandexRealtyScraper,
|
||||||
)
|
)
|
||||||
|
from scraper_kit.proxy_errors import NoProxyAvailableError
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
from sqlalchemy.orm import Session
|
from sqlalchemy.orm import Session
|
||||||
|
|
@ -208,7 +210,11 @@ def ban_kind_of_exception(exc: BaseException) -> str:
|
||||||
Один узел на все avito-сайты mark_banned: разводить исход в каждом из трёх
|
Один узел на все avito-сайты mark_banned: разводить исход в каждом из трёх
|
||||||
было бы тремя копиями одного условия.
|
было бы тремя копиями одного условия.
|
||||||
"""
|
"""
|
||||||
if isinstance(exc, AvitoSidecarUnavailableError):
|
if isinstance(exc, AvitoSidecarUnavailableError | NoProxyAvailableError):
|
||||||
|
# NoProxyAvailableError (#2687): пул прокси пуст в проде. Причина здесь не
|
||||||
|
# «не установлена» (тогда был бы 'unknown', #2764), а названа самим типом —
|
||||||
|
# см. proxy_errors.py: «это НАША инфраструктура (нет живого прокси), не
|
||||||
|
# внешний блок». Тот же диагноз, что и у отказа сайдкара.
|
||||||
return BAN_KIND_INFRA
|
return BAN_KIND_INFRA
|
||||||
if isinstance(exc, AvitoBlockedError | AvitoRateLimitedError):
|
if isinstance(exc, AvitoBlockedError | AvitoRateLimitedError):
|
||||||
return BAN_KIND_PLATFORM
|
return BAN_KIND_PLATFORM
|
||||||
|
|
@ -3217,6 +3223,20 @@ async def run_cian_full_load(
|
||||||
)
|
)
|
||||||
return counters
|
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:
|
except RuntimeError as exc:
|
||||||
# on_bucket кидает RuntimeError("cancelled") при cooperative cancel,
|
# on_bucket кидает RuntimeError("cancelled") при cooperative cancel,
|
||||||
# RuntimeError("shutdown") при кооперативном SIGTERM-drain (#1182 Phase 3a).
|
# RuntimeError("shutdown") при кооперативном SIGTERM-drain (#1182 Phase 3a).
|
||||||
|
|
@ -3445,6 +3465,20 @@ async def run_yandex_full_load(
|
||||||
)
|
)
|
||||||
return counters
|
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:
|
except RuntimeError as exc:
|
||||||
# on_bucket кидает RuntimeError("cancelled") при cooperative cancel,
|
# on_bucket кидает RuntimeError("cancelled") при cooperative cancel,
|
||||||
# RuntimeError("shutdown") при кооперативном SIGTERM-drain (#1182 Phase 3a).
|
# RuntimeError("shutdown") при кооперативном SIGTERM-drain (#1182 Phase 3a).
|
||||||
|
|
@ -3662,6 +3696,23 @@ async def run_avito_full_load(
|
||||||
)
|
)
|
||||||
return counters
|
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:
|
except RuntimeError as exc:
|
||||||
# on_bucket кидает RuntimeError("cancelled") при cooperative cancel,
|
# on_bucket кидает RuntimeError("cancelled") при cooperative cancel,
|
||||||
# RuntimeError("shutdown") при кооперативном SIGTERM-drain (#1182 Phase 3a).
|
# RuntimeError("shutdown") при кооперативном SIGTERM-drain (#1182 Phase 3a).
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue