"""Чекпоинт по якорям для 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, "исправный якорь не зафиксирован"