"""Каждая запись counters domclick-свипа несёт чекпоинт `done_buckets` (#3369). Доводка #3355/PR #3363: дрейн- и бан-ветки `run_domclick_city_sweep` кладут `done_buckets` в payload явно, а cancel-ветка, финальный heartbeat и финализаторы писали голый `counters.to_dict()`. Это НЕ было потерей данных: writer — `scraper_kit.orchestration.runs`, все его писатели мержат jsonb (`counters = COALESCE(counters,'{}') || CAST(:counters AS jsonb)`), а `_pick_resume` наследует ключ при claim (#3074). Инвариант ставится ради единообразия и независимости веток от того, какой писатель окажется на другом конце (перезаписывающий двойник `app/services/scrape_runs.py` существует). Проверка по ЗНАЧЕНИЮ, а не по факту записи: чекпоинт обязан дословно лежать в payload'е отмены, финального heartbeat'а и `mark_done`. """ from __future__ import annotations import os # Settings собирается автофикстурой conftest'а и требует database_url — как в # test_3355_drain_mark_full_loads.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; SELECT по `rid` отдаёт чекпоинт.""" def __init__(self, prev_counters: dict[str, Any] | None = None) -> None: self.writes: list[tuple[str, dict[str, Any]]] = [] self._prev_counters = prev_counters 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: row = types.SimpleNamespace(counters=self._prev_counters) return MagicMock(**{"fetchone.return_value": row if self._prev_counters else None}) return MagicMock() def commit(self) -> None: ... def rollback(self) -> None: ... 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: raise AssertionError(f"cancel-ветка не сработала: тронут скрапер ({name})") 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_domclick_sweep_cancel_keeps_inherited_checkpoint() -> None: """Отмена до SERP-фазы: payload несёт ровно унаследованные две корзины.""" 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: True), ): await pl.run_domclick_city_sweep( db, # type: ignore[arg-type] run_id=3369, config=_config(), matcher=MagicMock(), request_delay_sec=0.0, shutdown_requested=lambda: False, resume_run_id=3368, ) assert db.writes, "cancel не оставил ни одной записи counters" _sql, counters = db.writes[-1] assert counters.get("done_buckets") == inherited, ( f"cancel: payload без чекпоинта ({counters.get('done_buckets')!r} вместо " f"{inherited!r}) — ветка зависит от того, мержит ли писатель jsonb" ) @pytest.mark.asyncio async def test_domclick_sweep_cancel_marker_is_not_resume_status_only() -> None: """Отмена без чекпоинта в предшественнике не выдумывает корзин (пустой список).""" from scraper_kit.orchestration import pipeline as pl db = _FakeDb() with ( patch.object(pl, "DomClickScraper", _NeverCalledScraper), patch.object(pl.runs, "is_cancelled", lambda *_a: True), ): await pl.run_domclick_city_sweep( db, # type: ignore[arg-type] run_id=3369, config=_config(), matcher=MagicMock(), request_delay_sec=0.0, shutdown_requested=lambda: False, ) assert db.writes, "cancel не оставил ни одной записи counters" _sql, counters = db.writes[-1] assert counters.get("done_buckets") == [], ( f"cancel без предшественника: ожидался пустой чекпоинт, получено " f"{counters.get('done_buckets')!r}" ) class _CleanSweepScraper: """Свип прошёл все корзины без блока: ветка честного статуса → mark_done.""" def __init__(self, *_a: Any, **_kw: Any) -> None: self.blocked = False self.geo_filtered = 0 self.fetch_errors = 0 self.buckets_completed = 1 self.buckets_total = 1 self.completed_buckets = ["3:0:5000000"] self.bucket_start_index = 0 self.request_delay_sec = 0.0 async def __aenter__(self) -> _CleanSweepScraper: return self async def __aexit__(self, *_e: Any) -> None: return None async def fetch_city(self, **_kw: Any) -> list[Any]: return [] @pytest.mark.asyncio async def test_domclick_sweep_final_heartbeat_and_mark_done_carry_checkpoint() -> None: """Финальный heartbeat и mark_done несут унаследованное ∪ пройденное этим прогоном.""" from scraper_kit.orchestration import pipeline as pl inherited = ["1:0:5000000", "2:0:5000000"] expected = sorted([*inherited, "3:0:5000000"]) db = _FakeDb(prev_counters={"done_buckets": inherited}) with ( patch.object(pl, "DomClickScraper", _CleanSweepScraper), patch.object(pl.runs, "is_cancelled", lambda *_a: False), ): await pl.run_domclick_city_sweep( db, # type: ignore[arg-type] run_id=3369, config=_config(), matcher=MagicMock(), request_delay_sec=0.0, shutdown_requested=lambda: False, resume_run_id=3368, ) assert len(db.writes) >= 2, f"ожидались heartbeat + финализатор, есть {len(db.writes)}" sql_done, counters_done = db.writes[-1] assert "finished_at" in sql_done, f"последняя запись — не финализатор: {sql_done!r}" _sql_hb, counters_hb = db.writes[-2] for label, payload in (("финальный heartbeat", counters_hb), ("mark_done", counters_done)): assert payload.get("done_buckets") == expected, ( f"{label}: payload без чекпоинта ({payload.get('done_buckets')!r} вместо " f"{expected!r}) — ветка зависит от того, мержит ли писатель jsonb" )