diff --git a/tradein-mvp/backend/tests/test_3074_avito_anchor_checkpoint.py b/tradein-mvp/backend/tests/test_3074_avito_anchor_checkpoint.py new file mode 100644 index 00000000..f7c0b488 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3074_avito_anchor_checkpoint.py @@ -0,0 +1,175 @@ +"""Чекпоинт по якорям для avito_city_sweep (#3074). + +Прод-факт, из которого выросла задача. Типичный итог свипа: + + {"anchors_done": 1, "anchors_total": 5, ..., "enrichment_abort_note": + "detail enrichment aborted (Avito detail firewall/soft-block ...)"} + +То есть прогон срывается блокировкой на ПЕРВОМ из пяти якорей — 25 банов за +60 дней. Без чекпоинта следующий прогон снова начинает с первого якоря, +упирается в ту же стену, и якоря 2-5 не собираются никогда. + +Ключ чекпоинта — ИМЯ якоря, а не его индекс: состав списка зависит от +`city_slug`, и позиция в нём не устойчива между городами. + +ИНВАРИАНТ, РАДИ КОТОРОГО ТЕСТ. В чекпоинт попадает только якорь, пройденный до +конца. Ветка блокировки делает `return` и до записи не доходит, а вот +`except Exception` — доходит: якорь упал, но цикл продолжается. Записать такой +якорь как пройденный значило бы, что следующий прогон его пропустит и +объявления оттуда не соберутся НИКОГДА, причём молча — прогон завершится +штатно. Ровно та же граница, что у combo в yandex-свипе. +""" + +from __future__ import annotations + +import os + +# Settings собирается автофикстурой conftest'а и требует database_url. Выставляем +# до остальных импортов — так же, как в test_3074_yandex_sweep_checkpoint.py. +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +import json +import types +from typing import Any +from unittest.mock import MagicMock, patch + +import pytest + +ANCHOR_A = (56.83, 60.60, "ekb-center") +ANCHOR_B = (56.79, 60.63, "ekb-south") + + +class _FakeDb: + """Резюм-SELECT отдаёт counters предшественника; heartbeat'ы записываются.""" + + def __init__(self, prev_counters: dict[str, Any] | None = None) -> None: + self.prev_counters = prev_counters or {} + self.heartbeats: list[dict[str, Any]] = [] + + def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any: + if params and "counters" in params: + self.heartbeats.append(json.loads(params["counters"])) + return MagicMock() + if params and "rid" in params: + return MagicMock( + fetchone=lambda: types.SimpleNamespace(counters=self.prev_counters) + ) + return MagicMock() + + def commit(self) -> None: ... + def rollback(self) -> None: ... + + +class _FakeAsyncSession: + def __init__(self, *_a: Any, **_kw: Any) -> None: ... + async def __aenter__(self) -> _FakeAsyncSession: + return self + + async def __aexit__(self, *_e: Any) -> None: + return None + + +class _FakeScraper: + """Двойник AvitoScraper: помнит, за какими якорями реально ходили.""" + + visited: list[tuple[float, float]] = [] # noqa: RUF012 — тестовый сборник + raise_on: tuple[float, float] | None = None + + def __init__(self, *_a: Any, **_kw: Any) -> None: + self._browser = None + self._cffi = None + + async def fetch_around(self, lat: float, lon: float, *_a: Any, **_kw: Any) -> list: + _FakeScraper.visited.append((lat, lon)) + if _FakeScraper.raise_on == (lat, lon): + raise RuntimeError("якорь упал по не-баново́й причине") + return [] + + +def _config() -> types.SimpleNamespace: + return types.SimpleNamespace( + scraper_fetch_mode="cffi", + scraper_proxy_url=None, + use_proxy_pool_browser=False, + browser_http_endpoint=None, + environment="test", + avito_serp_ok_not_banned=True, + ) + + +async def _run(prev: dict[str, Any] | None, raise_on: tuple[float, float] | None = None): + from scraper_kit.orchestration import pipeline as pl + + _FakeScraper.visited = [] + _FakeScraper.raise_on = raise_on + db = _FakeDb(prev) + + with ( + patch.object(pl, "AvitoScraper", _FakeScraper), + patch.object(pl, "AsyncSession", _FakeAsyncSession), + patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_avito_city_sweep( + db, # type: ignore[arg-type] + run_id=7001, + config=_config(), + matcher=MagicMock(), + enrichment=MagicMock(), + anchors=[ANCHOR_A, ANCHOR_B], + enrich_houses=False, + enrich_imv=False, + detail_top_n=0, + resume_run_id=6999 if prev is not None else None, + ) + return db + + +def _last_checkpoint(db: _FakeDb) -> list[str]: + with_ckpt = [hb for hb in db.heartbeats if "done_buckets" in hb] + assert with_ckpt, "ни один heartbeat не унёс done_buckets — чекпоинт не персистится" + return with_ckpt[-1]["done_buckets"] + + +@pytest.mark.asyncio +async def test_checkpointed_anchor_is_skipped_without_a_single_request() -> None: + """Якорь из чекпоинта не опрашивается вовсе — ни одного обращения к источнику. + + Ядро задачи: до фикса повторный прогон снова шёл в первый якорь и снова + получал там бан. + """ + await _run({"done_buckets": ["ekb-center"]}) + + assert (ANCHOR_A[0], ANCHOR_A[1]) not in _FakeScraper.visited, ( + "якорь из чекпоинта всё-таки опрашивали" + ) + assert (ANCHOR_B[0], ANCHOR_B[1]) in _FakeScraper.visited, "второй якорь не обошли" + + +@pytest.mark.asyncio +async def test_checkpoint_accumulates_over_inherited() -> None: + """Пройденный якорь дописывается поверх унаследованных, а не затирает их.""" + db = await _run({"done_buckets": ["ekb-center"]}) + assert _last_checkpoint(db) == ["ekb-center", "ekb-south"] + + +@pytest.mark.asyncio +async def test_without_resume_all_anchors_are_visited() -> None: + """Без чекпоинта поведение прежнее — обходятся все якоря.""" + db = await _run(None) + assert len(_FakeScraper.visited) == 2 + assert _last_checkpoint(db) == ["ekb-center", "ekb-south"] + + +@pytest.mark.asyncio +async def test_failed_anchor_does_not_enter_checkpoint() -> None: + """Упавший якорь НЕ считается пройденным. + + Иначе следующий прогон пропустит его навсегда, и это будет незаметно: + прогон завершается штатно, просто часть города не собирается никогда. + """ + db = await _run(None, raise_on=(ANCHOR_A[0], ANCHOR_A[1])) + + ckpt = _last_checkpoint(db) + assert "ekb-center" not in ckpt, "упавший якорь попал в чекпоинт" + assert "ekb-south" in ckpt, "исправный якорь не зафиксирован" 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 c8520c19..9d0c9bdd 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 @@ -1132,6 +1132,7 @@ async def run_avito_city_sweep( request_delay_sec: float = 7.0, enrich_imv: bool = True, region_code: int = DEFAULT_REGION_CODE, + resume_run_id: int | None = None, ) -> CitySweepCounters: """Full city sweep: iterate anchors × pages → save → enrich houses + detail → IMV. @@ -1154,6 +1155,32 @@ async def run_avito_city_sweep( # а avito_slug (#12) — в путь URL, скоупя сам запрос на город-цель вместо ЕКБ. # None → ЕКБ-дефолт (совпадает с anchors=EKB_ANCHORS fallback ниже). _anchors = anchors if anchors is not None else EKB_ANCHORS + + # ── Checkpoint/resume (#3074): единица обхода — ЯКОРЬ ────────────────────── + # Прод-факт: у avito_city_sweep типичный итог `anchors_done: 1` из + # `anchors_total: 5` — прогон срывается блокировкой детализации на первом же + # якоре (25 банов за 60 дней). Без чекпоинта следующий прогон снова начинает + # с первого якоря, упирается в ту же стену, и якоря 2-5 не собираются никогда. + # + # Ключ чекпоинта — ИМЯ якоря, а не индекс: список якорей зависит от + # `city_slug`, и позиция в нём не устойчива между прогонами разных городов. + _skip_anchors: set[str] = set() + if resume_run_id is not None: + _prev = db.execute( + text("SELECT counters FROM scrape_runs WHERE id = CAST(:rid AS bigint)"), + {"rid": resume_run_id}, + ).fetchone() + if _prev is not None and _prev.counters: + _pc: dict = _prev.counters if isinstance(_prev.counters, dict) else {} + _skip_anchors = set(_pc.get("done_buckets", [])) + logger.info( + "city-sweep run_id=%d: resuming from run %s — %d якорей уже пройдено", + run_id, + resume_run_id, + len(_skip_anchors), + ) + _done_anchors: set[str] = set(_skip_anchors) + _loc = get_city_location(city_slug) # #262 wave 2: avito_slug у CityLocation Optional — не у каждого известного города # он подтверждён (403/429 на исчерпанном пуле при проверке, либо omonym-коллизия). @@ -1237,6 +1264,19 @@ async def run_avito_city_sweep( ) try: for idx, (lat, lon, name) in enumerate(_anchors, start=1): + if name in _skip_anchors: + # #3074: якорь собран предыдущим оборванным прогоном — ни одного + # HTTP-запроса. Счётчик двигаем, чтобы `anchors_done` продолжал + # означать «докуда дошли по списку», а не «сколько собрал этот run». + counters.anchors_done = idx + logger.info( + "city-sweep run_id=%d: anchor #%d/%d (%s) пропущен — есть в чекпоинте", + run_id, + idx, + len(_anchors), + name, + ) + continue if runs.is_cancelled(db, run_id): logger.info( "city-sweep run_id=%d: cancelled at anchor #%d/%d (%s)", @@ -1273,6 +1313,10 @@ async def run_avito_city_sweep( lon, ) + # #3074: сбрасывается на КАЖДОЙ итерации — иначе один упавший якорь + # заразил бы все последующие, и чекпоинт не пополнялся бы вовсе. + _anchor_ok = True + # Capture loop variables in default args (B023): prevents stale binding # if the coroutine is scheduled after the loop variable changes. _a_lat, _a_lon, _a_name = lat, lon, name @@ -1767,9 +1811,20 @@ async def run_avito_city_sweep( except Exception: logger.exception("city-sweep run_id=%d: anchor %s failed", run_id, name) counters.errors_count += 1 + _anchor_ok = False counters.anchors_done = idx - runs.update_heartbeat(db, run_id, counters.to_dict()) + # #3074: в чекпоинт попадает ТОЛЬКО якорь, пройденный до конца. + # Ветка блокировки выше делает `return` и сюда не доходит, а вот + # generic-except доходит — якорь упал, но цикл продолжается. Записать + # его как пройденный значило бы, что следующий прогон его пропустит и + # объявления оттуда не соберутся НИКОГДА, причём молча: прогон + # завершится штатно. Тот же инвариант, что у combo в yandex-свипе. + if _anchor_ok: + _done_anchors.add(name) + runs.update_heartbeat( + db, run_id, {**counters.to_dict(), "done_buckets": sorted(_done_anchors)} + ) # ── IMV-фаза: финальный обход тронутых домов ────────── if enrich_imv and all_touched_house_ids: diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py index 3a65a109..2a384cdf 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py @@ -786,6 +786,10 @@ async def _job_avito_city_sweep( request_delay_sec=float(params.get("request_delay_sec", 7.0)), enrich_houses=bool(params.get("enrich_houses", True)), radius_m=int(params.get("radius_m", 1500)), + # #3074: подхват якорей у оборванного предшественника. Ключ чекпоинта — + # имя якоря, оно не зависит от количества якорей, поэтому в отличие от + # combo-чекпоинта yandex-свипа гарда по числу якорей здесь не требуется. + resume_run_id=_pick_resume(db, run_id), )