From e12ece8fe6181dc704c6a0804da05db466d711e1 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sun, 6 Sep 2026 09:45:56 +0500 Subject: [PATCH 1/2] =?UTF-8?q?fix(tradein/cian):=20houses-=D1=84=D0=B0?= =?UTF-8?q?=D0=B7=D0=B0=20cian=5Fcity=5Fsweep=20=D0=B8=D0=B4=D1=91=D1=82?= =?UTF-8?q?=20=D1=87=D0=B5=D1=80=D0=B5=D0=B7=20=D0=BF=D1=83=D0=BB=20=D0=BF?= =?UTF-8?q?=D1=80=D0=BE=D0=BA=D1=81=D0=B8=20+=20=D1=81=D1=82=D0=BE=D0=BF?= =?UTF-8?q?=20=D0=BD=D0=B0=20=D0=BF=D1=83=D1=81=D1=82=D0=BE=D0=BC=20=D0=BF?= =?UTF-8?q?=D1=83=D0=BB=D0=B5=20(#3394)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `run_cian_city_sweep` звал `fetch_newbuilding(zhk_url, config=config)` без `proxy_provider`, хотя провайдер лежит в аргументах самого свипа и соседние фазы (SERP через CianScraper, detail через cian_fetch_detail) его передают. Внутри это давало `build_browser_fetcher(config, "cian", proxy_provider=None)` → `use_pool` эффективно False → POST /fetch без "proxy" → сайдкар брал свой env-узел SCRAPER_PROXY_URL, на проде выключенный (407 → camoufox InvalidIP → /fetch 503, факт #3386). Фаза houses давала 0/30 с 02.09 (run 6179: 24 × `houses failed ... 503 Service Unavailable`), строк `override=True` в логах сайдкара по ней не было ни одной. Остальные вызывающие `fetch_newbuilding` / `resolve_cian_zhk_url_via_search` провайдер уже передают (#2767/#2830/#3382), этот вызов был последним мимо пула. Второе: пустой пул поднимается ДО запроса, следующий дом упрётся ровно в то же самое — общий `except Exception` на дом превращал это в 30 одинаковых houses_failed и прогон уходил в 'done'. Теперь NoProxyAvailableError рвёт фазу и свип: `no_proxy_stop=1` в counters, mark_banned с ban_kind='infra' и сохранённым done_buckets (образец — #3389 yandex-nb-sweep, #3382). Ветка стоит ДО generic-except, иначе наш отказ инфраструктуры читался бы как «IP likely blocked» — бан площадки. --- .../test_3394_cian_sweep_houses_proxy_pool.py | 222 ++++++++++++++++++ .../src/scraper_kit/orchestration/pipeline.py | 65 ++++- 2 files changed, 285 insertions(+), 2 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_3394_cian_sweep_houses_proxy_pool.py diff --git a/tradein-mvp/backend/tests/test_3394_cian_sweep_houses_proxy_pool.py b/tradein-mvp/backend/tests/test_3394_cian_sweep_houses_proxy_pool.py new file mode 100644 index 00000000..e80e1d4f --- /dev/null +++ b/tradein-mvp/backend/tests/test_3394_cian_sweep_houses_proxy_pool.py @@ -0,0 +1,222 @@ +"""#3394 — houses-фаза `cian_city_sweep` ходила в сайдкар мимо прокси-пула. + +`run_cian_city_sweep` звал `fetch_newbuilding(zhk_url, config=config)` БЕЗ +`proxy_provider`, хотя провайдер лежал прямо в аргументах свипа и SERP/detail-фазы +рядом его передают. Внутри `fetch_newbuilding` это давало +`build_browser_fetcher(config, "cian", proxy_provider=None)` → `use_pool` эффективно +False → POST /fetch без "proxy" → сайдкар брал env-узел `SCRAPER_PROXY_URL`, на проде +выключенный (407 → camoufox `InvalidIP` → `/fetch` 503, факт #3386). Фаза давала 0/30 с +02.09 (run 6179: 24 × `houses failed ... 503 Service Unavailable`). + +Тесты идут через НАСТОЯЩИЙ `run_cian_city_sweep` и НАСТОЯЩИЙ `fetch_newbuilding`; +подделан только сам `BrowserFetcher` в `providers/_base` — то есть проверяется тот +тракт, по которому строится фетчер на проде, а не отдельный вызов из теста. + +Второй тест — про исход пустого пула: он поднимается ДО запроса, следующий дом упрётся +ровно в то же самое, поэтому фаза (и свип) обязаны оборваться на первом доме с +`no_proxy_stop=1`, а не писать 30 одинаковых отказов и уходить в 'done' (образец — +#3389, yandex-nb-sweep). + +Сеть/БД/камуфокс замоканы; в сеть тест не ходит. +""" + +from __future__ import annotations + +import os +from types import SimpleNamespace +from typing import Any, ClassVar +from unittest.mock import AsyncMock, MagicMock, patch + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +import pytest +from scraper_kit.orchestration.pipeline import run_cian_city_sweep +from scraper_kit.proxy_errors import NoProxyAvailableError + +PFX = "scraper_kit.orchestration.pipeline" +_BASE_FETCHER = "scraper_kit.providers._base.BrowserFetcher" + +_ZHK_URL = "https://zhk-test-ekb-i.cian.ru/" + + +class _CapturingFetcher: + """Собирает kwargs конструктора; `fetch()` отдаёт HTML без initialState.""" + + captured: ClassVar[list[dict[str, Any]]] = [] + + def __init__(self, **kwargs: Any) -> None: + _CapturingFetcher.captured.append(kwargs) + + async def __aenter__(self) -> _CapturingFetcher: + return self + + async def __aexit__(self, *_: object) -> None: + return None + + async def fetch(self, _url: str, **_kwargs: Any) -> str: + return "" + + def report_ban(self, _reason: str) -> None: + return None + + +class _EmptyPoolFetcher: + """Пул пуст: `fetch()` поднимает `NoProxyAvailableError` ДО запроса (#2616).""" + + calls: ClassVar[list[str]] = [] + + def __init__(self, **_kwargs: Any) -> None: + pass + + async def __aenter__(self) -> _EmptyPoolFetcher: + return self + + async def __aexit__(self, *_: object) -> None: + return None + + async def fetch(self, url: str, **_kwargs: Any) -> str: + _EmptyPoolFetcher.calls.append(url) + raise NoProxyAvailableError("cian") + + def report_ban(self, _reason: str) -> None: + return None + + +class _RunsRecorder: + """Двойник scrape_runs: пишет финализаторы, is_cancelled всегда False.""" + + def __init__(self) -> None: + self.calls: list[tuple[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: + self.calls.append(("update_heartbeat", dict(counters))) + + def mark_done(self, _db: Any, _run_id: int, counters: dict[str, Any]) -> None: + self.calls.append(("mark_done", dict(counters))) + + def mark_banned( + self, + _db: Any, + _run_id: int, + _error: str, + counters: dict[str, Any], + *, + ban_kind: str = "unknown", + ) -> None: + self.calls.append(("mark_banned", 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 _config() -> SimpleNamespace: + return SimpleNamespace( + scraper_fetch_mode="curl_cffi", + browser_http_endpoint="http://tradein-browser:3000/fetch", + scraper_proxy_url=None, + cian_proxy_url=None, + use_proxy_pool_browser=True, + environment="production", + cian_proxy_max_rotations=0, + proxy_rotate_attempts=1, + proxy_rotate_attempt_timeout_s=1.0, + scraper_skip_seen_today=False, + ) + + +def _db(house_rows: list[dict[str, Any]]) -> MagicMock: + """db.execute: houses-запрос отдаёт дома, любой другой — пусто.""" + + def _execute(stmt: Any, *_a: Any, **_kw: Any) -> MagicMock: + res = MagicMock() + rows = house_rows if "FROM houses h" in str(stmt) else [] + res.mappings.return_value.all.return_value = rows + res.fetchone.return_value = None + res.scalar_one.return_value = 0 + return res + + db = MagicMock() + db.execute.side_effect = _execute + return db + + +def _nb_lot() -> MagicMock: + return MagicMock( + listing_segment="novostroyki", + house_source="cian_newbuilding", + house_ext_id="48853", + ) + + +async def _drive( + *, fetcher: type, house_rows: list[dict[str, Any]] +) -> tuple[dict[str, int], list[tuple[str, dict[str, Any]]], MagicMock]: + recorder = _RunsRecorder() + provider = MagicMock(name="proxy_provider") + scraper = MagicMock() + scraper.__aenter__ = AsyncMock(return_value=scraper) + scraper.__aexit__ = AsyncMock(return_value=None) + scraper.fetch_around_multi_room = AsyncMock(return_value=[_nb_lot()]) + scraper.state_extraction_attempts = 1 + scraper.state_extraction_failures = 0 + with ( + patch(f"{PFX}.CianScraper", return_value=scraper), + patch(f"{PFX}.save_listings", MagicMock(return_value=(1, 0))), + patch(f"{PFX}.runs", recorder), + patch(_BASE_FETCHER, fetcher), + ): + counters = await run_cian_city_sweep( + _db(house_rows), + config=_config(), + matcher=MagicMock(), + run_id=1, + proxy_provider=provider, + anchors=[(56.84, 60.60, "A1")], + radius_m=1000, + pages_per_anchor=1, + request_delay_sec=0.0, + detail_top_n=0, + enrich_houses=True, + newbuilding_only=True, + ) + return counters.to_dict(), recorder.calls, provider + + +@pytest.mark.asyncio +async def test_houses_phase_builds_fetcher_with_proxy_provider() -> None: + """(а) фетчер houses-фазы собран С провайдером пула — тот, что пришёл в свип.""" + _CapturingFetcher.captured.clear() + counters, _calls, provider = await _drive( + fetcher=_CapturingFetcher, + house_rows=[{"id": 7, "cian_zhk_url": _ZHK_URL}], + ) + assert counters["houses_attempted"] == 1 + assert len(_CapturingFetcher.captured) == 1, _CapturingFetcher.captured + kwargs = _CapturingFetcher.captured[0] + # Ровно те три аргумента, без которых сайдкар уходит на env-узел. + assert kwargs["proxy_provider"] is provider + assert kwargs["use_pool"] is True + assert kwargs["environment"] == "production" + + +@pytest.mark.asyncio +async def test_empty_pool_stops_houses_phase_after_first_house() -> None: + """(б) пустой пул → один дом, no_proxy_stop=1, свип оборван (mark_banned).""" + _EmptyPoolFetcher.calls.clear() + counters, calls, _provider = await _drive( + fetcher=_EmptyPoolFetcher, + house_rows=[ + {"id": 7, "cian_zhk_url": _ZHK_URL}, + {"id": 8, "cian_zhk_url": _ZHK_URL}, + {"id": 9, "cian_zhk_url": _ZHK_URL}, + ], + ) + assert _EmptyPoolFetcher.calls == [_ZHK_URL], "второй дом упрётся в то же самое" + assert counters["houses_attempted"] == 1 + assert counters["houses_failed"] == 1 + assert counters["no_proxy_stop"] == 1 + assert calls[-1][0] == "mark_banned" + assert calls[-1][1]["no_proxy_stop"] == 1 diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py index 0f41f6b4..ada518f9 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py @@ -78,7 +78,7 @@ from scraper_kit.providers.yandex.serp import ( ROOM_PATH, YandexRealtyScraper, ) -from scraper_kit.proxy_errors import NoProxyAvailableError +from scraper_kit.proxy_errors import NoProxyAvailableError, caused_by_no_proxy if TYPE_CHECKING: from sqlalchemy.orm import Session @@ -2845,6 +2845,11 @@ class CianCitySweepCounters: houses_enriched: int = 0 houses_failed: int = 0 errors_count: int = 0 + # #3394: 1 — прогон оборван «пул прокси пуст» (к площадке не ходили), 0 — нет. + # int, а не bool: счётчики читают SQL'ем `counters->>'no_proxy_stop' = '1'`, + # мимо JSON-ного `true` он промахнётся молча (тот же довод — в + # app/tasks/yandex_newbuilding_sweep.py::SweepResult.to_dict). + no_proxy_stop: int = 0 def to_dict(self) -> dict[str, int]: return {f.name: getattr(self, f.name) for f in fields(self)} @@ -3221,7 +3226,18 @@ async def run_cian_city_sweep( zhk_url: str = hrow["cian_zhk_url"] counters.houses_attempted += 1 try: - nb_enrichment = await fetch_newbuilding(zhk_url, config=config) + # proxy_provider (#3394): без него build_browser_fetcher внутри + # fetch_newbuilding получал proxy_provider=None → use_pool + # эффективно False → POST /fetch уходил БЕЗ "proxy", и сайдкар + # брал свой env-узел SCRAPER_PROXY_URL — на проде выключенный + # (407 → camoufox InvalidIP → /fetch 503). Фаза houses шла 0/30 + # с 02.09, хотя detail-фаза выше ходила пулом всё это время. + # Соседние вызывающие fetch_newbuilding провайдер передают: + # cian_history_backfill.py:352 (#3382), + # newbuilding_enrich_backfill.py:546 (#2767), admin.py:2175. + nb_enrichment = await fetch_newbuilding( + zhk_url, config=config, proxy_provider=proxy_provider + ) if nb_enrichment is not None: save_newbuilding_enrichment(db, house_id_val, nb_enrichment) counters.houses_enriched += 1 @@ -3237,6 +3253,23 @@ async def run_cian_city_sweep( except Exception as exc: counters.houses_failed += 1 counters.errors_count += 1 + if caused_by_no_proxy(exc): + # #3394 (образец — #3197/#3389): пустой пул это НЕ отказ + # площадки — запрос не уходил вовсе, и следующий дом + # упрётся ровно в то же самое. Перебирать оставшиеся + # значило бы писать 30 одинаковых houses_failed и уходить + # в 'done'. Обрываем фазу и свип (ветка + # `except NoProxyAvailableError` в цикле по anchor'ам). + logger.error( + "cian-sweep run_id=%d: houses СТОП — пул прокси пуст, " + "к площадке не ходили (house_id=%d, дом %d/%d): %s", + run_id, + house_id_val, + hidx + 1, + len(house_rows), + exc, + ) + raise logger.warning( "cian-sweep run_id=%d: houses failed house_id=%d: %s", run_id, @@ -3298,6 +3331,34 @@ async def run_cian_city_sweep( counters.anchors_done = idx runs.update_heartbeat(db, run_id, counters.to_dict()) continue + except NoProxyAvailableError as exc: + # #3394: пул пуст — свип обрывается на первом же доме/лоте, а не + # перебирает остаток города одинаковыми отказами. Ветка стоит ДО + # generic-except: там NoProxyAvailableError уходил бы в + # consecutive_failures и (через 3 якоря) в mark_banned с диагнозом + # «IP likely blocked», то есть наш отказ читался бы как бан площадки. + # done_buckets сохраняем — следующий прогон не начнёт с нуля (#2686). + counters.no_proxy_stop = 1 + counters.errors_count += 1 + logger.error( + "cian-sweep run_id=%d: СТОП на anchor #%d/%d (%s) — пул прокси " + "пуст, к площадке не ходили: %s", + run_id, + idx, + len(_anchors), + name, + exc, + ) + _ckpt = {**counters.to_dict(), "done_buckets": sorted(_done_anchors)} + runs.update_heartbeat(db, run_id, _ckpt) + runs.mark_banned( + db, + run_id, + f"cian sweep aborted: {exc}", + _ckpt, + ban_kind=ban_kind_of_exception(exc), + ) + return counters except Exception as e: logger.exception("cian-sweep run_id=%d: anchor %s SERP failed", run_id, name) counters.errors_count += 1 From fdae823762352825dba8cf19df48935485a6a8f1 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sun, 6 Sep 2026 10:15:43 +0500 Subject: [PATCH 2/2] =?UTF-8?q?fix(#3394):=20=D1=81=D1=82=D0=BE=D0=BF=20?= =?UTF-8?q?=D0=BF=D0=BE=20=D0=BF=D1=83=D1=81=D1=82=D0=BE=D0=BC=D1=83=20?= =?UTF-8?q?=D0=BF=D1=83=D0=BB=D1=83=20=E2=80=94=20=D0=BF=D0=BE=20=D1=86?= =?UTF-8?q?=D0=B5=D0=BF=D0=BE=D1=87=D0=BA=D0=B5=20=D0=BF=D1=80=D0=B8=D1=87?= =?UTF-8?q?=D0=B8=D0=BD=20=D0=BD=D0=B0=20=D1=83=D1=80=D0=BE=D0=B2=D0=BD?= =?UTF-8?q?=D0=B5=20=D1=8F=D0=BA=D0=BE=D1=80=D1=8F;=20=D0=BE=D0=B4=D0=B8?= =?UTF-8?q?=D0=BD=20errors=5Fcount;=20ban=5Fkind=3Dinfra=20=D0=BF=D0=BE?= =?UTF-8?q?=D0=B4=20=D1=82=D0=B5=D1=81=D1=82=D0=BE=D0=BC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../test_3394_cian_sweep_houses_proxy_pool.py | 84 ++++++++++++++-- .../src/scraper_kit/orchestration/pipeline.py | 96 ++++++++++++------- 2 files changed, 137 insertions(+), 43 deletions(-) diff --git a/tradein-mvp/backend/tests/test_3394_cian_sweep_houses_proxy_pool.py b/tradein-mvp/backend/tests/test_3394_cian_sweep_houses_proxy_pool.py index e80e1d4f..1169e218 100644 --- a/tradein-mvp/backend/tests/test_3394_cian_sweep_houses_proxy_pool.py +++ b/tradein-mvp/backend/tests/test_3394_cian_sweep_houses_proxy_pool.py @@ -14,8 +14,13 @@ False → POST /fetch без "proxy" → сайдкар брал env-узел `S Второй тест — про исход пустого пула: он поднимается ДО запроса, следующий дом упрётся ровно в то же самое, поэтому фаза (и свип) обязаны оборваться на первом доме с -`no_proxy_stop=1`, а не писать 30 одинаковых отказов и уходить в 'done' (образец — -#3389, yandex-nb-sweep). +`no_proxy_stop=1` и `ban_kind='infra'`, а не писать 30 одинаковых отказов и уходить в +'done' (образец — #3389, yandex-nb-sweep). + +Третий — тот же исход, когда провайдер ЗАВЕРНУЛ отказ пула в своё исключение: ловля на +уровне якоря идёт по цепочке причин (`caused_by_no_proxy`), а не по типу, иначе первая же +такая обёртка (у `fetch_detail` она уже есть) увела бы стоп в generic-обработчик и +записала бы нашу же нехватку прокси баном площадки. Сеть/БД/камуфокс замоканы; в сеть тест не ходит. """ @@ -82,11 +87,37 @@ class _EmptyPoolFetcher: return None +class _WrappingEmptyPoolFetcher: + """Пул пуст, но провайдер завернул отказ в своё исключение (как `fetch_detail`).""" + + calls: ClassVar[list[str]] = [] + + def __init__(self, **_kwargs: Any) -> None: + pass + + async def __aenter__(self) -> _WrappingEmptyPoolFetcher: + return self + + async def __aexit__(self, *_: object) -> None: + return None + + async def fetch(self, url: str, **_kwargs: Any) -> str: + _WrappingEmptyPoolFetcher.calls.append(url) + raise RuntimeError("wrapped") from NoProxyAvailableError("cian") + + def report_ban(self, _reason: str) -> None: + return None + + class _RunsRecorder: """Двойник scrape_runs: пишет финализаторы, is_cancelled всегда False.""" def __init__(self) -> None: self.calls: list[tuple[str, dict[str, Any]]] = [] + # ban_kind писался в scrape_runs, но двойник его отбрасывал — поле уходило + # из-под теста целиком, а именно оно отличает «наша инфраструктура» ('infra') + # от «площадка забанила» ('platform'/'unknown'). + self.ban_kinds: list[str] = [] def is_cancelled(self, _db: Any, _run_id: int) -> bool: return False @@ -107,6 +138,7 @@ class _RunsRecorder: ban_kind: str = "unknown", ) -> None: self.calls.append(("mark_banned", dict(counters))) + self.ban_kinds.append(ban_kind) def mark_failed(self, _db: Any, _run_id: int, _error: str, counters: dict[str, Any]) -> None: self.calls.append(("mark_failed", dict(counters))) @@ -153,7 +185,7 @@ def _nb_lot() -> MagicMock: async def _drive( *, fetcher: type, house_rows: list[dict[str, Any]] -) -> tuple[dict[str, int], list[tuple[str, dict[str, Any]]], MagicMock]: +) -> tuple[dict[str, int], _RunsRecorder, MagicMock]: recorder = _RunsRecorder() provider = MagicMock(name="proxy_provider") scraper = MagicMock() @@ -182,21 +214,25 @@ async def _drive( enrich_houses=True, newbuilding_only=True, ) - return counters.to_dict(), recorder.calls, provider + return counters.to_dict(), recorder, provider @pytest.mark.asyncio async def test_houses_phase_builds_fetcher_with_proxy_provider() -> None: """(а) фетчер houses-фазы собран С провайдером пула — тот, что пришёл в свип.""" _CapturingFetcher.captured.clear() - counters, _calls, provider = await _drive( + counters, _recorder, provider = await _drive( fetcher=_CapturingFetcher, house_rows=[{"id": 7, "cian_zhk_url": _ZHK_URL}], ) assert counters["houses_attempted"] == 1 assert len(_CapturingFetcher.captured) == 1, _CapturingFetcher.captured kwargs = _CapturingFetcher.captured[0] - # Ровно те три аргумента, без которых сайдкар уходит на env-узел. + # Эпохи различает РОВНО первый assert: `use_pool`/`environment` + # build_browser_fetcher кладёт из config независимо от провайдера (_base.py:233), + # то есть и до правки они были такими же. Оставлены как контракт сайдкара: без + # любого из трёх (или при use_pool=False / environment != production) запрос + # уходит на env-узел SCRAPER_PROXY_URL. assert kwargs["proxy_provider"] is provider assert kwargs["use_pool"] is True assert kwargs["environment"] == "production" @@ -206,7 +242,7 @@ async def test_houses_phase_builds_fetcher_with_proxy_provider() -> None: async def test_empty_pool_stops_houses_phase_after_first_house() -> None: """(б) пустой пул → один дом, no_proxy_stop=1, свип оборван (mark_banned).""" _EmptyPoolFetcher.calls.clear() - counters, calls, _provider = await _drive( + counters, recorder, _provider = await _drive( fetcher=_EmptyPoolFetcher, house_rows=[ {"id": 7, "cian_zhk_url": _ZHK_URL}, @@ -218,5 +254,35 @@ async def test_empty_pool_stops_houses_phase_after_first_house() -> None: assert counters["houses_attempted"] == 1 assert counters["houses_failed"] == 1 assert counters["no_proxy_stop"] == 1 - assert calls[-1][0] == "mark_banned" - assert calls[-1][1]["no_proxy_stop"] == 1 + assert counters["errors_count"] == 1, "одно событие — один errors_count" + assert recorder.calls[-1][0] == "mark_banned" + assert recorder.calls[-1][1]["no_proxy_stop"] == 1 + assert recorder.ban_kinds[-1] == "infra", "наша инфраструктура, не бан площадки" + + +@pytest.mark.asyncio +async def test_empty_pool_stops_even_when_wrapped_by_provider() -> None: + """(в) отказ пула, ЗАВЁРНУТЫЙ провайдером, тоже обрывает свип — по цепочке причин. + + Решение «стоп» принимает raise-site (`caused_by_no_proxy`, тип в цепочке + `__cause__`), а ловля стояла по конкретному типу `NoProxyAvailableError`. Пока + `fetch_newbuilding` не оборачивает, это совпадало; `fetch_detail` уже оборачивает + (см. app/tasks/cian_history_backfill.py:236-238), и с появлением такой обёртки стоп + уходил бы в generic-обработчик → consecutive_failures → mark_banned «IP likely + blocked», то есть НАШ отказ инфраструктуры записывался бы баном площадки. + """ + _WrappingEmptyPoolFetcher.calls.clear() + counters, recorder, _provider = await _drive( + fetcher=_WrappingEmptyPoolFetcher, + house_rows=[ + {"id": 7, "cian_zhk_url": _ZHK_URL}, + {"id": 8, "cian_zhk_url": _ZHK_URL}, + {"id": 9, "cian_zhk_url": _ZHK_URL}, + ], + ) + assert _WrappingEmptyPoolFetcher.calls == [_ZHK_URL], "второй дом упрётся в то же самое" + assert counters["houses_attempted"] == 1 + assert counters["no_proxy_stop"] == 1 + assert counters["errors_count"] == 1 + assert recorder.calls[-1][0] == "mark_banned" + assert recorder.ban_kinds[-1] == "infra" diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py index ada518f9..2f7e204c 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py @@ -2846,9 +2846,12 @@ class CianCitySweepCounters: houses_failed: int = 0 errors_count: int = 0 # #3394: 1 — прогон оборван «пул прокси пуст» (к площадке не ходили), 0 — нет. - # int, а не bool: счётчики читают SQL'ем `counters->>'no_proxy_stop' = '1'`, - # мимо JSON-ного `true` он промахнётся молча (тот же довод — в - # app/tasks/yandex_newbuilding_sweep.py::SweepResult.to_dict). + # int, а не bool — как у всех соседей (avito_detail_backfill.py:1043, + # domclick_detail_backfill.py:610, yandex_newbuilding_sweep.py::SweepResult.to_dict): + # весь payload counters числовой, и SQL-разбор вида `counters->>'no_proxy_stop' = '1'` + # на нём не промахнётся, тогда как мимо JSON-ного `true` промахнулся бы молча. + # Такого запроса в репо пока НЕТ (грепом 06.09 — ни в data/sql, ни в коде) — это + # конвенция формата на будущее, а не поддержка существующего потребителя. no_proxy_stop: int = 0 def to_dict(self) -> dict[str, int]: @@ -3127,6 +3130,23 @@ async def run_cian_city_sweep( _cian_detail_consec_failures = 0 except Exception as exc: counters.detail_failed += 1 + if caused_by_no_proxy(exc): + # #3394: тот же довод, что в houses-фазе ниже — пул пуст, + # запрос не уходил, следующий лот упрётся ровно в то же + # самое. Без этой ветки отказ доедал _cian_detail_abort и + # писал в лог «proxy likely banned/stuck» (диагноз + # площадки на НАШЕЙ инфраструктуре), а перед этим ещё + # дёргал ротацию IP, которой пустому пулу нечего менять. + # errors_count инкрементируется один раз — на уровне + # якоря, где прогон обрывается. + logger.error( + "cian-sweep run_id=%d: detail СТОП — пул прокси пуст, " + "к площадке не ходили (listing_id=%d): %s", + run_id, + listing_id, + exc, + ) + raise counters.errors_count += 1 _cian_detail_consec_failures += 1 logger.warning( @@ -3252,14 +3272,14 @@ async def run_cian_city_sweep( ) except Exception as exc: counters.houses_failed += 1 - counters.errors_count += 1 if caused_by_no_proxy(exc): # #3394 (образец — #3197/#3389): пустой пул это НЕ отказ # площадки — запрос не уходил вовсе, и следующий дом # упрётся ровно в то же самое. Перебирать оставшиеся # значило бы писать 30 одинаковых houses_failed и уходить - # в 'done'. Обрываем фазу и свип (ветка - # `except NoProxyAvailableError` в цикле по anchor'ам). + # в 'done'. Обрываем фазу и свип (ветка caused_by_no_proxy + # в цикле по anchor'ам). errors_count инкрементируется там + # же и только там — иначе одно событие считалось дважды. logger.error( "cian-sweep run_id=%d: houses СТОП — пул прокси пуст, " "к площадке не ходили (house_id=%d, дом %d/%d): %s", @@ -3270,6 +3290,7 @@ async def run_cian_city_sweep( exc, ) raise + counters.errors_count += 1 logger.warning( "cian-sweep run_id=%d: houses failed house_id=%d: %s", run_id, @@ -3331,35 +3352,42 @@ async def run_cian_city_sweep( counters.anchors_done = idx runs.update_heartbeat(db, run_id, counters.to_dict()) continue - except NoProxyAvailableError as exc: - # #3394: пул пуст — свип обрывается на первом же доме/лоте, а не - # перебирает остаток города одинаковыми отказами. Ветка стоит ДО - # generic-except: там NoProxyAvailableError уходил бы в - # consecutive_failures и (через 3 якоря) в mark_banned с диагнозом - # «IP likely blocked», то есть наш отказ читался бы как бан площадки. - # done_buckets сохраняем — следующий прогон не начнёт с нуля (#2686). - counters.no_proxy_stop = 1 - counters.errors_count += 1 - logger.error( - "cian-sweep run_id=%d: СТОП на anchor #%d/%d (%s) — пул прокси " - "пуст, к площадке не ходили: %s", - run_id, - idx, - len(_anchors), - name, - exc, - ) - _ckpt = {**counters.to_dict(), "done_buckets": sorted(_done_anchors)} - runs.update_heartbeat(db, run_id, _ckpt) - runs.mark_banned( - db, - run_id, - f"cian sweep aborted: {exc}", - _ckpt, - ban_kind=ban_kind_of_exception(exc), - ) - return counters except Exception as e: + if caused_by_no_proxy(e): + # #3394: пул пуст — свип обрывается на первом же доме/лоте, а не + # перебирает остаток города одинаковыми отказами. Ловим по ЦЕПОЧКЕ + # причин и ДО generic-разбора ниже: raise-site внутри фаз решает + # «стоп» тем же caused_by_no_proxy, и стоит провайдеру завернуть + # отказ в своё исключение (так уже делает fetch_detail — см. + # app/tasks/cian_history_backfill.py:236-238), как ветка по ТИПУ + # промахнулась бы: наш отказ инфраструктуры уехал бы в + # consecutive_failures и (через 3 якоря) в mark_banned с диагнозом + # «IP likely blocked», то есть читался бы как бан площадки. + # ban_kind — константой, а не ban_kind_of_exception(e): тот смотрит + # isinstance и на завёрнутом отказе дал бы 'unknown'; здесь диагноз + # установлен самим условием ветки. + # done_buckets сохраняем — следующий прогон не начнёт с нуля (#2686). + counters.no_proxy_stop = 1 + counters.errors_count += 1 + logger.error( + "cian-sweep run_id=%d: СТОП на anchor #%d/%d (%s) — пул прокси " + "пуст, к площадке не ходили: %s", + run_id, + idx, + len(_anchors), + name, + e, + ) + _ckpt = {**counters.to_dict(), "done_buckets": sorted(_done_anchors)} + runs.update_heartbeat(db, run_id, _ckpt) + runs.mark_banned( + db, + run_id, + f"cian sweep aborted: {e}", + _ckpt, + ban_kind=BAN_KIND_INFRA, + ) + return counters logger.exception("cian-sweep run_id=%d: anchor %s SERP failed", run_id, name) counters.errors_count += 1 consecutive_failures += 1