"""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 # Данные-атрибуты объявляем явно: catch-all __getattr__ ниже отдаёт корутину # на ЛЮБОЕ имя, и счётчик #3368 (читается в finally, ревью #3373) уехал бы # в counters функцией — падало бы сериализацией, а не смыслом. self.capped_buckets = 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 (#3390) тут не спасает: он сливает payload со СВОЕЙ строкой прогона, а та создана пустой (`create_run`) — унаследованный чекпоинт лежит в counters ПРЕДЫДУЩЕГО прогона, и scheduler при claim его не переносит. Без явного ключа дрейн-прогон закрывался бы с пустым чекпоинтом, и резюм пересобирал бы все шесть корзин заново. """ 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)})" )