diff --git a/tradein-mvp/backend/tests/test_3074_cian_anchor_checkpoint.py b/tradein-mvp/backend/tests/test_3074_cian_anchor_checkpoint.py new file mode 100644 index 00000000..0a2d75f4 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3074_cian_anchor_checkpoint.py @@ -0,0 +1,159 @@ +"""Чекпоинт по якорям для cian_city_sweep (#3074). + +Замер за 60 дней, по которому выбран источник: 65 прогонов, среднее 35 минут, +максимум 72, две отмены деплоем. Пятиминутного дренажа (#3029) на такие прогоны +не хватает — убитый на 35-й минуте сбор начинался заново с первого якоря. + +Ключ чекпоинта — ИМЯ якоря, а не индекс: состав списка зависит от `city_slug` +(областные свипы идут по своим наборам), позиция между городами не устойчива. + +ИНВАРИАНТ. В чекпоинт попадает только якорь, пройденный до конца. У циана эта +граница уже проведена потоком управления: все ветки отказа делают `return` или +`continue`, и до записи не доходят. Этим он отличается от avito-свипа, где успех +и неудача сходились в одной строке и потребовался отдельный флаг. Тест ниже +проверяет, что граница не нарушена: упавший якорь не должен попасть в чекпоинт, +иначе следующий прогон пропустит его навсегда — молча, потому что прогон +завершится штатно, просто часть города не соберётся. +""" + +from __future__ import annotations + +import os + +# Settings собирается автофикстурой conftest'а и требует database_url. +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: + 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 _FakeScraper: + """Двойник CianScraper: помнит, за какими якорями реально ходили.""" + + visited: list[tuple[float, float]] = [] # noqa: RUF012 — тестовый сборник + raise_on: tuple[float, float] | None = None + + def __init__(self, *_a: Any, **_kw: Any) -> None: + # Счётчики, которые конвейер читает у скрапера после каждого якоря + # (#2625 — диагностика «все якоря упали»). Без них падает не проверяемая + # логика, а сам двойник. + self.state_extraction_attempts = 1 + self.state_extraction_failures = 0 + self.request_delay_sec = 0.0 + + async def __aenter__(self) -> _FakeScraper: + return self + + async def __aexit__(self, *_e: Any) -> None: + return None + + async def fetch_around_multi_room(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_proxy_url=None, + scraper_fetch_mode="cffi", + use_proxy_pool_browser=False, + browser_http_endpoint=None, + environment="test", + ) + + +async def _run(prev: dict[str, Any] | None, raise_on: tuple[float, float] | None = None) -> _FakeDb: + from scraper_kit.orchestration import pipeline as pl + + _FakeScraper.visited = [] + _FakeScraper.raise_on = raise_on + db = _FakeDb(prev) + + with ( + patch.object(pl, "CianScraper", _FakeScraper), + patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_cian_city_sweep( + db, # type: ignore[arg-type] + run_id=8001, + config=_config(), + matcher=MagicMock(), + anchors=[ANCHOR_A, ANCHOR_B], + enrich_houses=False, + detail_top_n=0, + request_delay_sec=0.0, + resume_run_id=7999 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: + """Упавший якорь НЕ считается пройденным. + + У циана это обеспечено потоком управления, а не флагом: ветка отказа делает + `continue` до записи. Тест сторожит именно это — рефакторинг, сливающий ветки + в одну, сломает инвариант незаметно. + """ + 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 9d0c9bdd..b743ebc6 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 @@ -2766,6 +2766,7 @@ async def run_cian_city_sweep( enrich_houses: bool = True, newbuilding_only: bool = True, region_code: int = DEFAULT_REGION_CODE, + resume_run_id: int | None = None, ) -> CianCitySweepCounters: """Cian newbuilding city sweep: SERP → detail(+price-history) → newbuilding/houses. @@ -2782,6 +2783,33 @@ async def run_cian_city_sweep( mark_done вызывается ВСЕГДА (finally outer). """ _anchors = anchors if anchors is not None else EKB_ANCHORS + + # ── Checkpoint/resume (#3074): единица обхода — ЯКОРЬ ────────────────────── + # Замер за 60 дней: 65 прогонов, среднее 35 минут, максимум 72, две отмены + # деплоем. Пятиминутного дренажа (#3029) на такие прогоны не хватает, и + # убитый на 35-й минуте сбор начинается заново с первого якоря. + # + # Ключ — ИМЯ якоря, а не индекс: состав списка зависит от `city_slug` + # (областные свипы идут по своим наборам), позиция между городами не + # устойчива. По той же причине гарда по числу якорей не нужна — в отличие от + # combo-чекпоинта яндекса, где ключ якоря не содержал. + _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( + "cian-sweep run_id=%d: resuming from run %s — %d якорей уже пройдено", + run_id, + resume_run_id, + len(_skip_anchors), + ) + _done_anchors: set[str] = set(_skip_anchors) + # city_slug (#12): region_id города-цели → CianScraper.city_region_id скоупит SERP # на город вместо дефолтного ЕКБ. None/неизвестный slug → ЕКБ-дефолт в конструкторе. _loc = get_city_location(city_slug) @@ -2827,6 +2855,19 @@ async def run_cian_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( + "cian-sweep run_id=%d: anchor #%d/%d (%s) пропущен — есть в чекпоинте", + run_id, + idx, + len(_anchors), + name, + ) + continue # Cooperative cancel перед каждым anchor if runs.is_cancelled(db, run_id): logger.info( @@ -3178,7 +3219,19 @@ async def run_cian_city_sweep( continue counters.anchors_done = idx - runs.update_heartbeat(db, run_id, counters.to_dict()) + # #3074: сюда попадают ТОЛЬКО якоря, пройденные до конца. Все ветки + # отказа выше либо `return`, либо `continue` — в чекпоинт они не + # заходят. Разница с avito-свипом, где успех и неудача сходились в + # одной строке и потребовался отдельный флаг: здесь граница уже + # проведена самим потоком управления, и её достаточно не нарушать. + # + # Записать упавший якорь пройденным значило бы, что следующий прогон + # пропустит его навсегда — молча, потому что прогон завершится + # штатно, просто часть города не соберётся. + _done_anchors.add(name) + runs.update_heartbeat( + db, run_id, {**counters.to_dict(), "done_buckets": sorted(_done_anchors)} + ) # Пауза между anchor'ами (поверх per-request sleep внутри scraper'а) if idx < len(_anchors): 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 2a384cdf..9eaed46c 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 @@ -920,6 +920,10 @@ async def _job_cian_city_sweep( detail_top_n=int(params.get("detail_top_n", 10)), enrich_houses=bool(params.get("enrich_houses", True)), newbuilding_only=bool(params.get("newbuilding_only", True)), + # #3074: подхват якорей у оборванного предшественника. Ключ чекпоинта — + # имя якоря, от их количества не зависит, поэтому гарда как у combo- + # чекпоинта яндекса здесь не требуется. + resume_run_id=_pick_resume(db, run_id), )