diff --git a/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py b/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py index 569c7d1b..bb434bb7 100644 --- a/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py +++ b/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py @@ -45,12 +45,26 @@ import pytest class _FakeDb: """Каждый UPDATE с counters запоминается вместе с текстом SQL (для статуса).""" - def __init__(self) -> None: + def __init__( + self, + prev_counters: dict[str, Any] | None = None, + detail_rows: list[dict[str, Any]] | None = None, + ) -> None: self.writes: list[tuple[str, dict[str, Any]]] = [] + self._prev_counters = prev_counters + self._detail_rows = detail_rows or [] def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any: if params and "counters" in params: self.writes.append((str(stmt), json.loads(params["counters"]))) + return MagicMock() + if params and "rid" in params: + # SELECT counters предыдущего прогона (resume_run_id → чекпоинт). + row = types.SimpleNamespace(counters=self._prev_counters) + return MagicMock(**{"fetchone.return_value": row if self._prev_counters else None}) + if params and "lim" in params: + # SELECT кандидатов detail-обогащения (cian full-load). + return MagicMock(**{"mappings.return_value.all.return_value": self._detail_rows}) return MagicMock() def commit(self) -> None: ... @@ -225,6 +239,92 @@ async def test_domclick_sweep_drain_is_marked_interrupted() -> None: ) +@pytest.mark.asyncio +async def test_domclick_sweep_drain_keeps_inherited_checkpoint() -> None: + """domclick: дрейн ОБЯЗАН переписать унаследованный чекпоинт в свой payload. + + Мержа jsonb тут нет: app-level `scrape_runs.update_heartbeat`/`mark_done` пишут + `counters = CAST(:counters AS jsonb)` — полная перезапись, а scheduler при claim + done_buckets не наследует. Без явного ключа дрейн-прогон закрывался бы с пустым + чекпоинтом, и резюм пересобирал бы все шесть корзин заново. + """ + from scraper_kit.orchestration import pipeline as pl + + inherited = ["1:0:5000000", "2:0:5000000"] + db = _FakeDb(prev_counters={"done_buckets": inherited}) + with ( + patch.object(pl, "DomClickScraper", _NeverCalledScraper), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_domclick_city_sweep( + db, # type: ignore[arg-type] + run_id=3355, + config=_config(), + matcher=MagicMock(), + request_delay_sec=0.0, + shutdown_requested=lambda: True, + resume_run_id=3354, + ) + + assert len(db.writes) >= 2, "дрейн обязан писать и heartbeat, и финализатор" + for label, (_sql, counters) in (("heartbeat", db.writes[-2]), ("финализатор", db.writes[-1])): + assert counters.get("done_buckets") == inherited, ( + f"{label}: чекпоинт потерян ({counters.get('done_buckets')!r} вместо " + f"{inherited!r}) — резюм пересоберёт уже собранные корзины" + ) + assert counters.get("interrupted") == 1, f"{label}: дрейн без метки" + + +@pytest.mark.asyncio +async def test_cian_full_load_drain_in_detail_phase_is_marked_interrupted() -> None: + """cian full-load: дрейн в detail-фазе — SERP целый, но прогон НЕ полный. + + Ранний выход из detail-цикла делал `break` и уходил в обычный `mark_done`: + статус 'done' без метки, хотя обогащение обрезано на первой же записи. + """ + from scraper_kit.orchestration import pipeline as pl + + class _SerpDoneScraper: + """SERP отработал целиком (on_bucket не зовётся) — дрейн ловит detail-фаза.""" + + def __init__(self, *_a: Any, **_kw: Any) -> None: + self._browser = None + self.request_delay_sec = 0.0 + self.last_dropped_nb = 0 + self.state_extraction_attempts = 0 + self.state_extraction_failures = 0 + + async def __aenter__(self) -> _SerpDoneScraper: + return self + + async def __aexit__(self, *_e: Any) -> None: + return None + + async def fetch_all_secondary(self, **_kw: Any) -> None: + return None + + db = _FakeDb(detail_rows=[{"id": 1, "source_url": "https://cian.ru/1"}]) + with ( + patch.object(pl, "CianScraper", _SerpDoneScraper), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_cian_full_load( + db, # type: ignore[arg-type] + run_id=3355, + config=_config(), + matcher=MagicMock(), + request_delay_sec=0.0, + enrich_detail=True, + detail_top_n=5, + shutdown_requested=lambda: True, + ) + + _status, counters = _final(db) + assert counters.get("interrupted") == 1, ( + "cian full-load: дрейн в detail-фазе финализирован как полный прогон" + ) + + @pytest.mark.asyncio async def test_zero_drain_status_is_resumable() -> None: """Минор ревью #3346: какой СТАТУС даёт нулевой дрейн — и берёт ли его резюм. 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 9fed62f2..f96db16a 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 @@ -3648,7 +3648,12 @@ async def run_cian_full_load( didx + 1, len(priority_rows), ) - break + # #3355: не break. SERP-фаза тут ЦЕЛАЯ (done_buckets полные), + # обрезано только обогащение — но break уводил в обычный + # mark_done ниже, и оборванный прогон выглядел полным. Тот же + # sentinel, что в _on_bucket → handler `RuntimeError("shutdown")` + # ставит interrupted=1 и сохраняет done_buckets. + raise RuntimeError("shutdown") listing_id: int = row["id"] source_url: str = row["source_url"] counters.detail_attempted += 1 @@ -4400,39 +4405,12 @@ async def run_domclick_city_sweep( _scraper_ref: list[DomClickScraper] = [] try: - # ── Cooperative cancel перед SERP-фазой ────────────────────────────── - if runs.is_cancelled(db, run_id): - logger.info("domclick-sweep run_id=%d: cancelled before SERP", run_id) - runs.update_heartbeat(db, run_id, counters.to_dict()) - return counters - elif shutdown_requested(): - # SIGTERM-drain до SERP-фазы: финализируем run как mark_done(partial), - # минуя honest-status (это чистый drain, не QRATOR-блок). - logger.info("domclick-sweep run_id=%d: SIGTERM-drain — stopping before SERP", run_id) - # #3355: метка дрейна (см. city-свипы #3330/#3333). done_buckets тут не - # пишем — унаследованный при claim (#3074) чекпоинт лежит в counters, и - # jsonb-мерж heartbeat/финализатора его сохраняет. - _drain = {**counters.to_dict(), "interrupted": 1} - runs.update_heartbeat(db, run_id, _drain) - runs.mark_done(db, run_id, _drain) - return counters - - logger.info( - "domclick-sweep run_id=%d: BFF citywide sweep city_id=%d " - "buckets=%d pages_cap=%d (watchdog %ds)", - run_id, - city_id, - len(ROOM_BUCKETS), - pages, - _sweep_timeout, - ) - - lots: list[ScrapedLot] = [] - # ── Checkpoint/resume (#3118): корзина ROOM_BUCKETS — единица обхода ── # Вместе со сдвигом #2854 подхват превращает случайную ротацию в # систематический обход: банимый на 1-й корзине источник закрывает все # шесть за несколько прогонов вместо повторов случайных. + # Читаем ДО ранних выходов: дрейн-ветка ниже обязана положить чекпоинт в + # свой payload сама (#3355), см. комментарий там. skip_buckets: set[str] = set() if resume_run_id is not None: _prev_row = db.execute( @@ -4451,6 +4429,44 @@ async def run_domclick_city_sweep( len(skip_buckets), ) + # ── Cooperative cancel перед SERP-фазой ────────────────────────────── + if runs.is_cancelled(db, run_id): + logger.info("domclick-sweep run_id=%d: cancelled before SERP", run_id) + runs.update_heartbeat(db, run_id, counters.to_dict()) + return counters + elif shutdown_requested(): + # SIGTERM-drain до SERP-фазы: финализируем run как mark_done(partial), + # минуя honest-status (это чистый drain, не QRATOR-блок). + logger.info("domclick-sweep run_id=%d: SIGTERM-drain — stopping before SERP", run_id) + # #3355: метка дрейна (см. city-свипы #3330/#3333) И явный перенос + # чекпоинта. Никакого jsonb-мержа тут нет: domclick пишет counters + # через app-level scrape_runs (update_heartbeat/mark_done делают + # `counters = CAST(:counters AS jsonb)` — ПОЛНАЯ перезапись), а + # scheduler при claim done_buckets не наследует. Без явного + # done_buckets дрейн-прогон закрылся бы с пустым чекпоинтом и резюм + # пересобирал бы все шесть корзин заново. Тот же приём, что в + # ban-ветке ниже (except NoProxyAvailableError). + _drain = { + **counters.to_dict(), + "done_buckets": sorted(skip_buckets), + "interrupted": 1, + } + runs.update_heartbeat(db, run_id, _drain) + runs.mark_done(db, run_id, _drain) + return counters + + logger.info( + "domclick-sweep run_id=%d: BFF citywide sweep city_id=%d " + "buckets=%d pages_cap=%d (watchdog %ds)", + run_id, + city_id, + len(ROOM_BUCKETS), + pages, + _sweep_timeout, + ) + + lots: list[ScrapedLot] = [] + async def _domclick_phase() -> None: """Единственная citywide-фаза: fetch_city + save.""" nonlocal lots