diff --git a/tradein-mvp/backend/tests/test_3074_avito_newbuilding_checkpoint.py b/tradein-mvp/backend/tests/test_3074_avito_newbuilding_checkpoint.py new file mode 100644 index 00000000..51f87df8 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3074_avito_newbuilding_checkpoint.py @@ -0,0 +1,266 @@ +"""Чекпоинт по страницам для avito_newbuilding_sweep (#3074). + +Последний длинный свип без чекпоинтов: 27.08 прогон убит деплоем на 343-й +минуте, собранное потеряно целиком — save_listings был ОДИН на весь sweep, +ждать в конце было нечего. Единица возобновления — СТРАНИЦА выдачи (у функции +уже есть `pages`; цикл по страницам живёт в `AvitoScraper._paginate_sweep`). + +Ключ чекпоинта — номер страницы. В отличие от якорей/combo, страницы строго +последовательны (break-on-empty), поэтому resume — это `start_page = +max(done_buckets) + 1`, без skip-набора произвольных элементов. + +ИНВАРИАНТЫ, РАДИ КОТОРЫХ ТЕСТ: +1. Резюм пропускает уже собранные страницы БЕЗ единого запроса к источнику. +2. В чекпоинт попадает только страница, пройденная до конца — оборванная + исключением/break страница НЕ фиксируется, иначе следующий прогон + пропустил бы её навсегда, причём молча (прогон завершится штатно). +3. Запись идёт мержем (`counters || :counters`), а не заменой — посторонний + ключ в `counters` текущего run'а переживает запись чекпоинта. +""" + +from __future__ import annotations + +import os + +# Settings собирается автофикстурой conftest'а и требует database_url. Выставляем +# до остальных импортов — так же, как в test_3074_avito_anchor_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 + + +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 _MergingFakeDb(_FakeDb): + """Как _FakeDb, но heartbeat-запись мержит в текущую строку (jsonb `||`), + а не просто копится списком — для теста инварианта #3 (посторонний ключ).""" + + def __init__(self, prev_counters: dict[str, Any] | None = None) -> None: + super().__init__(prev_counters) + self.row: dict[str, Any] = {} + + def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any: + if params and "counters" in params: + incoming = json.loads(params["counters"]) + self.row = {**self.row, **incoming} + self.heartbeats.append(dict(self.row)) + return MagicMock() + return super().execute(_stmt, params) + + +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 + + +def _lot(tag: str) -> MagicMock: + m = MagicMock() + m.source_id = tag + return m + + +class _FakeScraper: + """Двойник AvitoScraper: помнит реально пройденные страницы, умеет упасть.""" + + visited_pages: list[int] = [] # noqa: RUF012 — тестовый сборник + fail_on: int | None = None + + def __init__(self, *_a: Any, **_kw: Any) -> None: + self._browser = None + self._cffi = None + + async def fetch_newbuildings( + self, + *, + pages: int, + start_page: int = 1, + on_page: Any = None, + delay_override_sec: float | None = None, + ) -> list[Any]: + all_lots: list[Any] = [] + for page in range(start_page, pages + 1): + _FakeScraper.visited_pages.append(page) + if _FakeScraper.fail_on == page: + # Страница оборвана (сеть/парсинг) — break, on_page НЕ зовём. + break + new_lots = [_lot(f"p{page}-{i}") for i in range(2)] + all_lots.extend(new_lots) + if on_page is not None: + on_page(page, new_lots) + return all_lots + + +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", + ) + + +def _make_save(fail_on_call: int | None): + """save_listings, падающий на N-м вызове — имитация отказа БД на одной странице.""" + calls = {"n": 0} + + def _save(*_a, **_kw): + calls["n"] += 1 + if fail_on_call is not None and calls["n"] == fail_on_call: + raise RuntimeError(f"save_listings упал на странице {calls['n']}") + return (2, 0) + + return _save + + +async def _run( + db: _FakeDb, + *, + resume_run_id: int | None, + pages: int = 4, + fail_on: int | None = None, + save_fail_on_call: int | None = None, +): + from scraper_kit.orchestration import pipeline as pl + + _FakeScraper.visited_pages = [] + _FakeScraper.fail_on = fail_on + + with ( + patch.object(pl, "AvitoScraper", _FakeScraper), + patch.object(pl, "AsyncSession", _FakeAsyncSession), + patch.object(pl, "save_listings", _make_save(save_fail_on_call)), + ): + await pl.run_avito_newbuilding_sweep( + db, # type: ignore[arg-type] + run_id=8001, + config=_config(), + matcher=MagicMock(), + pages=pages, + request_delay_sec=0.0, + resume_run_id=resume_run_id, + ) + return db + + +def _last_checkpoint(db: _FakeDb) -> list[int]: + 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_resume_skips_done_pages_and_fetches_the_rest() -> None: + """Резюм с чекпоинтом [1,2] на pages=4 обходит только страницы 3 и 4 — + ни одного запроса к уже собранным.""" + db = _FakeDb({"done_buckets": [1, 2]}) + await _run(db, resume_run_id=7999, pages=4) + + assert _FakeScraper.visited_pages == [3, 4], ( + "уже собранные страницы 1/2 не должны опрашиваться повторно" + ) + assert _last_checkpoint(db) == [1, 2, 3, 4], "чекпоинт дописывается поверх унаследованного" + + +@pytest.mark.asyncio +async def test_without_resume_starts_from_page_one() -> None: + """Без resume_run_id поведение прежнее — обход с первой страницы.""" + db = _FakeDb(None) + await _run(db, resume_run_id=None, pages=3) + + assert _FakeScraper.visited_pages == [1, 2, 3] + assert _last_checkpoint(db) == [1, 2, 3] + + +@pytest.mark.asyncio +async def test_failed_page_does_not_enter_checkpoint() -> None: + """Страница, оборванная на середине (fail_on=2), НЕ считается пройденной. + + Иначе следующий прогон пропустит её навсегда молча — sweep завершится + штатно, просто часть выдачи не соберётся никогда. + """ + db = _FakeDb(None) + await _run(db, resume_run_id=None, pages=3, fail_on=2) + + # Страница 2 была АТАКОВАНА (попытка была — visited), но не пройдена до конца. + assert _FakeScraper.visited_pages == [1, 2] + ckpt = _last_checkpoint(db) + assert ckpt == [1], "упавшая страница 2 (и не начатая 3) не должны попасть в чекпоинт" + + +@pytest.mark.asyncio +async def test_checkpoint_write_merges_and_preserves_foreign_key() -> None: + """Запись чекпоинта — мерж (`counters || :counters`), не замена. + + Симулируем текущую строку run'а с посторонним ключом (например, оставленным + другим писателем/предыдущим heartbeat'ом) — он обязан пережить наши записи. + """ + db = _MergingFakeDb(None) + db.row = {"foreign_key": "survives-me"} + + await _run(db, resume_run_id=None, pages=1) + + assert db.row.get("foreign_key") == "survives-me", ( + "посторонний ключ в counters не пережил запись чекпоинта — запись была заменой, не мержем" + ) + assert "done_buckets" in db.row + + +@pytest.mark.asyncio +async def test_page_whose_save_failed_does_not_enter_checkpoint() -> None: + """Страница собрана, но её save упал — в чекпоинт она попасть не должна. + + Отказ save_listings перехватывается и прогон продолжается (это осознанно: + одна упавшая страница не должна ронять весь sweep). Но если отметить её + пройденной, следующий прогон её пропустит, и объявления оттуда не соберутся + НИКОГДА — молча, потому что прогон завершится штатно. + """ + db = _FakeDb(None) + await _run(db, resume_run_id=None, pages=3, save_fail_on_call=2) + + assert _FakeScraper.visited_pages == [1, 2, 3] + ckpt = _last_checkpoint(db) + assert 2 not in ckpt, "страница с упавшим save попала в чекпоинт — покрытие потеряно молча" + assert ckpt == [1, 3] + + +@pytest.mark.asyncio +async def test_resume_starts_at_first_gap_not_after_the_last_page() -> None: + """Дыра в чекпоинте перечитывается, а не перепрыгивается. + + `max(done)+1` пропустил бы страницу 2 навсегда. Продолжаем с первой + несобранной: страницы после дыры перечитаются, что дешевле потери и + безопасно — повторная запись схлопывается по dedup_hash. + """ + db = _FakeDb({"done_buckets": [1, 3]}) + await _run(db, resume_run_id=7777, pages=4) + + assert _FakeScraper.visited_pages[0] == 2, ( + "подхват начался не с дыры — пропущенная страница не соберётся никогда" + ) diff --git a/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py b/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py index ba629b56..d4e3eff7 100644 --- a/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py +++ b/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py @@ -495,10 +495,21 @@ async def _drive_nb_sweep( recorder = _RunsRecorder() db = MagicMock() lots = [MagicMock() for _ in range(6)] + + async def _fake_fetch_newbuildings( + *, pages: int, start_page: int = 1, on_page: Any = None, delay_override_sec: Any = None + ) -> list[Any]: + # #3074: реальный fetch_newbuildings зовёт on_page ПОСЛЕ каждой пройденной + # страницы (инкрементальный save) — двойник имитирует одну страницу с + # ВСЕМИ 6 лотами, чтобы save_mock/счётчики остались как раньше. + if on_page is not None: + on_page(start_page, lots) + return lots + scraper = MagicMock() scraper._cffi = None scraper._browser = None - scraper.fetch_newbuildings = AsyncMock(return_value=lots) + scraper.fetch_newbuildings = AsyncMock(side_effect=_fake_fetch_newbuildings) save_mock = MagicMock(side_effect=[(5, 1)]) avito_scraper_cls = MagicMock(return_value=scraper) if capture is not None: 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 99bd7952..97e70169 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 @@ -1947,6 +1947,7 @@ async def run_avito_newbuilding_sweep( pages: int = 20, request_delay_sec: float = 7.0, region_code: int = DEFAULT_REGION_CODE, + resume_run_id: int | None = None, ) -> NewbuildingSweepCounters: """Citywide-обход ЕКБ-выборки только новостроек (novostroyka-filter) → save. @@ -1962,9 +1963,91 @@ async def run_avito_newbuilding_sweep( Инжекция (#2135 F2): config/matcher/shutdown_requested приходят снаружи вместо прямых импортов app.* (см. scraper_kit.contracts). + + ── Checkpoint/resume (#3074): единица обхода — СТРАНИЦА выдачи ──────────── + Последний длинный свип без чекпоинтов: прогон убивался деплоем посреди + обхода (страница 300+ при pages в разы больше дефолта), и всё собранное + терялось целиком, потому что save_listings раньше был ОДИН на весь sweep — + ждать в конце было нечего. Здесь (как и у combo в yandex-свипе) save + происходит инкрементально в on_page-callback: сразу после того, как + страница пройдена до конца, её лоты сохраняются и её номер уходит в + чекпоинт `done_buckets` (список номеров страниц). Страницы — единственная и + строго последовательная единица обхода (break-on-empty), поэтому resume — + это просто `start_page = max(done_buckets) + 1`, без skip-набора: страница, + оборванная исключением, callback не получает и в чекпоинт не попадает — + иначе следующий прогон пропустил бы её навсегда, причём молча. """ counters = NewbuildingSweepCounters() + _done_pages: set[int] = 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 {} + _done_pages = {int(p) for p in _pc.get("done_buckets", [])} + logger.info( + "nb-sweep run_id=%d: resuming from run %s — %d страниц уже собрано", + run_id, + resume_run_id, + len(_done_pages), + ) + # #3074: продолжаем с ПЕРВОЙ несобранной страницы, а не с max+1. Дыра в + # чекпоинте возможна (страница собрана, но её save упал — она не отмечена, + # а следующие отмечены), и max+1 перепрыгнул бы её навсегда: молча, потому + # что прогон завершится штатно. Уже собранные страницы после дыры при этом + # перечитаются — это дешевле потери, а повторная запись идемпотентна + # (save_listings схлопывает по dedup_hash). + _start_page = 1 + while _start_page in _done_pages: + _start_page += 1 + + def _on_page(page: int, new_lots: list[ScrapedLot]) -> None: + """Страница собрана и СОХРАНЕНА — только тогда чекпоинт (мерж, не замена).""" + counters.lots_fetched += len(new_lots) + _saved_ok = True + if new_lots: + try: + # #2594: citywide novostroyka-обход — только ЕКБ (см. docstring). + ins, upd = save_listings( + db, + new_lots, + matcher=matcher, + region_code=region_code, + run_id=run_id, + city=EKATERINBURG_CITY_NAME, + ) + counters.lots_inserted += ins + counters.lots_updated += upd + except Exception as save_exc: + logger.exception( + "nb-sweep run_id=%d page=%d: save_listings failed: %s", + run_id, + page, + save_exc, + ) + counters.errors_count += 1 + _saved_ok = False + try: + db.rollback() + except Exception: + pass + # #3074: страница попадает в чекпоинт ТОЛЬКО если её лоты сохранены. + # Отказ save_listings перехвачен выше и прогон продолжается — но отметить + # страницу пройденной значило бы, что следующий прогон её пропустит, а + # объявления оттуда не соберутся НИКОГДА, причём молча: прогон завершится + # штатно. Тот же инвариант, что у якоря в avito city sweep и у бакета в + # cian full-load (_mark_bucket): в чекпоинт — только полностью собранная + # единица. Heartbeat обновляем в любом случае, иначе reap_zombies посчитает + # живой прогон мёртвым. + if _saved_ok: + _done_pages.add(page) + runs.update_heartbeat( + db, run_id, {**counters.to_dict(), "done_buckets": sorted(_done_pages)} + ) + browser_mode = config.scraper_fetch_mode == "browser" async with AsyncExitStack() as stack: session: AsyncSession | None = None @@ -2016,8 +2099,12 @@ async def run_avito_newbuilding_sweep( scraper._cffi = session try: - lots: list[ScrapedLot] = await scraper.fetch_newbuildings( + # #3074: save + чекпоинт идут инкрементально внутри _on_page — + # aggregate `lots` ниже используется только для лог-сообщения. + await scraper.fetch_newbuildings( pages=pages, + start_page=_start_page, + on_page=_on_page, delay_override_sec=request_delay_sec, ) except (AvitoBlockedError, AvitoRateLimitedError) as e: @@ -2027,30 +2114,6 @@ async def run_avito_newbuilding_sweep( ) return counters - counters.lots_fetched += len(lots) - if lots: - try: - # #2594: citywide novostroyka-обход — только ЕКБ (см. docstring). - ins, upd = save_listings( - db, - lots, - matcher=matcher, - region_code=region_code, - run_id=run_id, - city=EKATERINBURG_CITY_NAME, - ) - counters.lots_inserted += ins - counters.lots_updated += upd - except Exception as save_exc: - logger.exception( - "nb-sweep run_id=%d: save_listings failed: %s", run_id, save_exc - ) - counters.errors_count += 1 - try: - db.rollback() - except Exception: - pass - runs.update_heartbeat(db, run_id, counters.to_dict()) runs.mark_done(db, run_id, counters.to_dict()) logger.info( 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 3273a65f..0892d454 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 @@ -860,6 +860,9 @@ async def _job_avito_newbuilding_sweep( proxy_provider=ctx.proxy_provider, pages=int(params.get("pages", 20)), request_delay_sec=float(params.get("request_delay_sec", 7.0)), + # #3074: подхват страниц у оборванного предшественника — см. + # _job_avito_city_sweep выше. + resume_run_id=_pick_resume(db, run_id), ) diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/serp.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/serp.py index 5df6dbf8..bf017691 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/serp.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/serp.py @@ -298,6 +298,7 @@ def _avito_bisection_config(cap: int) -> BisectionConfig: probe_fail_policy=ProbeFailPolicy.SPLIT_OR_SKIP, ) + # HTTP 429 в curl_cffi-режиме через backconnect-прокси (mproxy.site) — НЕ IP-ban, а # transient «слишком много одновременных соединений» (лимит 5). Проходит на коротком # retry без ротации IP. Делаем до _AVITO_429_MAX_RETRIES коротких пауз; если они @@ -1803,13 +1804,28 @@ class AvitoScraper(BaseScraper): *, label: str, delay_override_sec: float | None = None, + start_page: int = 1, + on_page: Callable[[int, list[ScrapedLot]], None] | None = None, ) -> list[ScrapedLot]: """Общий paginated-обход ЕКБ (citywide / novostroyka) с break-on-empty. url_builder(page) — функция построения URL страницы (citywide или novostroyka). Сохраняет anti-block pipeline (_fetch_serp_html: firewall, IP rotation, ретраи), дедуп по source_id, break-on-empty. label — только - для логов. Не меняет наблюдаемое поведение fetch_city_wide. + для логов. start_page/on_page по умолчанию не меняют поведение — + fetch_city_wide вызывает без них. + + start_page (#3074): страница, с которой начать обход — страницы до + неё уже собраны предыдущим оборванным прогоном (checkpoint/resume), по + ним не делается ни одного HTTP-запроса. + on_page (#3074): опциональный callback(page: int, new_lots: list[ScrapedLot]) + -> None, вызывается СРАЗУ после того, как страница пройдена до конца + (в т.ч. с пустым new_lots — иначе последнюю пустую страницу + перечитывали бы вечно при resume). Страница, оборванная исключением + или break, callback не получает — вызов означает «страница пройдена + до конца». Позволяет инкрементальный save за пределами этого метода: + без него собранное часами обхода терялось целиком при убийстве + процесса (SIGKILL/деплой) — единственный save в конце ждать было нечем. """ if delay_override_sec is not None: self.request_delay_sec = delay_override_sec @@ -1817,7 +1833,7 @@ class AvitoScraper(BaseScraper): all_lots: list[ScrapedLot] = [] seen_ids: set[str] = set() - for page in range(1, pages + 1): + for page in range(start_page, pages + 1): url = url_builder(page) try: html = await self._fetch_serp_html_with_retry(url, page) @@ -1864,6 +1880,11 @@ class AvitoScraper(BaseScraper): len(lots) - len(new_lots), len(all_lots), ) + if on_page is not None: + try: + on_page(page, new_lots) + except Exception: + logger.exception("avito %s: on_page callback failed page=%d", label, page) if page < pages: await self.sleep_between_requests() return all_lots @@ -1873,6 +1894,8 @@ class AvitoScraper(BaseScraper): pages: int = 30, *, delay_override_sec: float | None = None, + start_page: int = 1, + on_page: Callable[[int, list[ScrapedLot]], None] | None = None, ) -> list[ScrapedLot]: """Обход ЕКБ-выборки только новостроек (novostroyka-filter), paginated. @@ -1886,6 +1909,10 @@ class AvitoScraper(BaseScraper): pages: максимальное число страниц (default 30). delay_override_sec: если задан — переопределяет request_delay_sec для этого вызова. + start_page (#3074): страница, с которой продолжить обход после + checkpoint/resume — см. _paginate_sweep. + on_page (#3074): callback(page, new_lots) после каждой пройденной + страницы — см. _paginate_sweep. Returns: Список ScrapedLot новостроек (все страницы, дедуп по source_id). @@ -1895,6 +1922,8 @@ class AvitoScraper(BaseScraper): self._build_newbuilding_url, label="newbuilding", delay_override_sec=delay_override_sec, + start_page=start_page, + on_page=on_page, ) logger.info("avito fetch_newbuildings pages=%d total_lots=%d", pages, len(all_lots)) return all_lots