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 new file mode 100644 index 00000000..bb434bb7 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py @@ -0,0 +1,369 @@ +"""SIGTERM-дрейн помечен `interrupted=1` и у full-load'ов с domclick-свипом (#3355). + +#3330/#3346 закрыли city-свипы, но четыре функции с ЖИВЫМ чекпоинтом `done_buckets` +продолжали финализировать дрейн чистым `done`: cian/yandex/avito full-load (ветка +`RuntimeError("shutdown")`) и domclick city sweep (дрейн до SERP-фазы). Резюм +(`_drained_done` в scheduler._resume_decision) такой прогон не берёт — чекпоинт +пишется в никуда. + +Кто чекпоинт ЧИТАЕТ (метка не диагностическая): + * cian_full_load — `skip_set` из `done_buckets` (ключ "room:lo:hi"), + планировщик зовёт с `resume_run_id=_pick_resume(...)`; + * avito_full_load — то же, ДВА job'а (обычный и exhaustive, #3315); + * domclick_city_sweep— `skip_buckets` из `done_buckets` (ROOM_BUCKETS), тоже + `_pick_resume`; + * yandex_full_load — читает `done_buckets` так же, НО у него нет job'а в + scheduler.py: единственный вход с `resume_run_id` — + админский POST. Метка там честно работает только при + ручном подхвате; автоматического потребителя нет. + +Статус финализации нулевого дрейна фиксируется тестом отдельно: honest-status +гейты (runs.mark_done → _sweep_run_did_nothing/_phase_totally_failed/ +_failed_ratio_too_high) могут увести прогон в 'failed'. Это резюмируемо +('failed' ∈ _RESUME_STATUSES), но статус — факт, а не догадка. + +Анти-цикл («чистый done не резюмируется») уже закрыт в +test_3333_drain_mark_all_sweeps.py — здесь не дублируется. +""" + +from __future__ import annotations + +import os + +# Settings собирается автофикстурой conftest'а и требует database_url — как в +# test_3333_drain_mark_all_sweeps.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: + """Каждый UPDATE с counters запоминается вместе с текстом SQL (для статуса).""" + + 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: ... + def rollback(self) -> None: ... + + +def _final(db: _FakeDb) -> tuple[str, dict[str, Any]]: + """Последняя запись counters = финализатор.""" + assert db.writes, "дрейн не оставил ни одной записи counters" + sql, counters = db.writes[-1] + if "status = 'done'" in sql: + return "done", counters + if "status = 'failed'" in sql: + return "failed", counters + return "heartbeat", counters + + +class _DrainAtFirstBucket: + """Любой fetch-метод сразу зовёт on_bucket — настоящий `_on_bucket` и роняет дрейн.""" + + def __init__(self, *_a: Any, **_kw: Any) -> None: + self._browser = None + self._cffi = None + self.request_delay_sec = 0.0 + + async def __aenter__(self) -> _DrainAtFirstBucket: + return self + + async def __aexit__(self, *_e: Any) -> None: + return None + + def __getattr__(self, name: str) -> Any: + async def _fetch(*_a: Any, **kw: Any) -> Any: + on_bucket = kw.get("on_bucket") + assert on_bucket is not None, f"{name} вызван без on_bucket — тест мимо дрейна" + on_bucket("2:0:5000000", []) + raise AssertionError("_on_bucket не оборвал прогон при shutdown_requested()") + + return _fetch + + +class _NeverCalledScraper: + """Дрейн срабатывает ДО скрапера: любой запрос тут — сломанный порядок проверок.""" + + def __init__(self, *_a: Any, **_kw: Any) -> None: ... + + async def __aenter__(self) -> _NeverCalledScraper: + return self + + async def __aexit__(self, *_e: Any) -> None: + return None + + def __getattr__(self, name: str) -> Any: + async def _boom(*_a: Any, **_kw: Any) -> Any: + raise AssertionError(f"скрапер вызван при дрейне: {name}") + + return _boom + + +def _config() -> types.SimpleNamespace: + return types.SimpleNamespace( + scraper_fetch_mode="cffi", + scraper_proxy_url=None, + scraper_skip_seen_today=False, + use_proxy_pool_browser=False, + browser_http_endpoint=None, + environment="test", + avito_serp_ok_not_banned=True, + cian_full_load_per_fetch_timeout_s=0.0, + ) + + +@pytest.mark.asyncio +async def test_cian_full_load_drain_is_marked_interrupted() -> None: + """cian full-load: дрейн на границе бакета → counters.interrupted == 1.""" + from scraper_kit.orchestration import pipeline as pl + + db = _FakeDb() + with ( + patch.object(pl, "CianScraper", _DrainAtFirstBucket), + 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, + shutdown_requested=lambda: True, + ) + + _status, counters = _final(db) + assert counters.get("interrupted") == 1, ( + "cian full-load: оборванный дрейном прогон неотличим от полного обхода — " + "резюм его чекпоинт не подхватит" + ) + + +@pytest.mark.asyncio +async def test_yandex_full_load_drain_is_marked_interrupted() -> None: + """yandex full-load: дрейн на границе бакета → counters.interrupted == 1.""" + from scraper_kit.orchestration import pipeline as pl + + db = _FakeDb() + with ( + patch.object(pl, "YandexRealtyScraper", _DrainAtFirstBucket), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_yandex_full_load( + db, # type: ignore[arg-type] + run_id=3355, + config=_config(), + matcher=MagicMock(), + enrichment=MagicMock(), + request_delay_sec=0.0, + shutdown_requested=lambda: True, + ) + + _status, counters = _final(db) + assert counters.get("interrupted") == 1, ( + "yandex full-load: оборванный дрейном прогон неотличим от полного обхода" + ) + + +@pytest.mark.asyncio +async def test_avito_full_load_drain_is_marked_interrupted() -> None: + """avito full-load (exhaustive, #3315): дрейн на границе бакета → interrupted == 1.""" + from scraper_kit.orchestration import pipeline as pl + + db = _FakeDb() + with ( + patch.object(pl, "AvitoScraper", _DrainAtFirstBucket), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_avito_full_load( + db, # type: ignore[arg-type] + run_id=3355, + config=_config(), + matcher=MagicMock(), + request_delay_sec=0.0, + shutdown_requested=lambda: True, + ) + + _status, counters = _final(db) + assert counters.get("interrupted") == 1, ( + "avito full-load: дрейн без метки — бан-бюджет тратится на уже собранные бакеты" + ) + + +@pytest.mark.asyncio +async def test_domclick_sweep_drain_is_marked_interrupted() -> None: + """domclick city sweep: дрейн до SERP-фазы → counters.interrupted == 1.""" + from scraper_kit.orchestration import pipeline as pl + + db = _FakeDb() + 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, + ) + + _status, counters = _final(db) + assert counters.get("interrupted") == 1, ( + "domclick: оборванный дрейном прогон неотличим от полного обхода" + ) + + +@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: какой СТАТУС даёт нулевой дрейн — и берёт ли его резюм. + + honest-status гейты в runs.mark_done (#2625/#2700/honest-run-status) вправе + увести прогон в 'failed'. Тест фиксирует ФАКТ, а не ожидание: обе ветки + резюмируемы ('failed' — через _RESUME_STATUSES, 'done' — через метку), но + третьего исхода быть не должно. + """ + from scraper_kit.orchestration import pipeline as pl + from scraper_kit.orchestration.scheduler import _RESUME_STATUSES, _resume_decision + + db = _FakeDb() + with ( + patch.object(pl, "AvitoScraper", _DrainAtFirstBucket), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + await pl.run_avito_full_load( + db, # type: ignore[arg-type] + run_id=3355, + config=_config(), + matcher=MagicMock(), + request_delay_sec=0.0, + shutdown_requested=lambda: True, + ) + + status, counters = _final(db) + assert status == "done", f"нулевой дрейн финализирован как {status!r} (было 'done')" + + prev = types.SimpleNamespace( + prev_id=4707, + prev_status=status, + prev_counters={**counters, "done_buckets": ["2:0:5000000"], "resume_chain": 0}, + same_params=True, + age_h=2.0, + interval_days="7", + ) + resume_from, verdict = _resume_decision(prev) + assert resume_from == 4707, ( + f"дрейн со статусом {status!r} не подхвачен резюмом: {verdict} " + f"(_RESUME_STATUSES={sorted(_RESUME_STATUSES)})" + ) 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 818acf74..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 @@ -3755,7 +3760,15 @@ async def run_cian_full_load( counters.saved_inserted, counters.saved_updated, ) - runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) + # #3355: та же метка дрейна, что у city-свипов (#3330/#3333). Без неё + # оборванный дрейном обход неотличим от полного: статус тот же 'done', + # а `_drained_done` в scheduler._resume_decision смотрит именно на неё — + # и done_buckets этого прогона не подхватывал никто. + runs.mark_done( + db, + run_id, + {**counters.to_dict(), "done_buckets": sorted(done), "interrupted": 1}, + ) return counters logger.exception("cian-full-load run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) @@ -3997,7 +4010,15 @@ async def run_yandex_full_load( counters.saved_inserted, counters.saved_updated, ) - runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) + # #3355: та же метка дрейна, что у city-свипов (#3330/#3333). Без неё + # оборванный дрейном обход неотличим от полного: статус тот же 'done', + # а `_drained_done` в scheduler._resume_decision смотрит именно на неё — + # и done_buckets этого прогона не подхватывал никто. + runs.mark_done( + db, + run_id, + {**counters.to_dict(), "done_buckets": sorted(done), "interrupted": 1}, + ) return counters logger.exception("yandex-full-load run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) @@ -4263,7 +4284,15 @@ async def run_avito_full_load( counters.saved_inserted, counters.saved_updated, ) - runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) + # #3355: та же метка дрейна, что у city-свипов (#3330/#3333). Без неё + # оборванный дрейном обход неотличим от полного: статус тот же 'done', + # а `_drained_done` в scheduler._resume_decision смотрит именно на неё — + # и done_buckets этого прогона не подхватывал никто. + runs.mark_done( + db, + run_id, + {**counters.to_dict(), "done_buckets": sorted(done), "interrupted": 1}, + ) return counters logger.exception("avito-full-load run_id=%d: fatal error", run_id) runs.mark_failed(db, run_id, str(exc), counters.to_dict()) @@ -4376,35 +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) - runs.update_heartbeat(db, run_id, counters.to_dict()) - runs.mark_done(db, run_id, counters.to_dict()) - 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( @@ -4423,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