"""Issue #3168: курсорные backfill-циклы не переживают рестарт (rosreestr_dkp_import). Разведка (см. PR): среди пяти backfill-циклов из issue (avito_detail_backfill, house_imv_backfill, cian_history_backfill, yandex_detail_backfill, geocode_missing_listings) НИ ОДИН не использует in-memory числовой курсор — все устроены как WHERE ... IS NULL / NOT EXISTS ... LIMIT (естественно резюмируемы без персистентной точки, см. cian_history_backfill.py и комментарий в scheduler.py про "Checkpoint/resume semantics"). Единственное место с буквально описанным в issue дефектом (курсор `WHERE id > last_id ORDER BY id`, последняя строчка heartbeat = "checkpoint", но last_id инициализируется литералом 0 при каждом запуске) — app/services/scheduler.py :: import_rosreestr_dkp (source='rosreestr_dkp_import'), шестой backfill, не входящий в таблицу issue. Фикс и тесты — про него; остальные пять на отдельный заход НЕ нужны, у них нет этого дефекта. Фикс — _resume_dkp_cursor(db, run_id) в scheduler.py: - Кандидат — последний прогон source='rosreestr_dkp_import' (кроме себя, кроме 'skipped'), тот же принцип, что и scraper_kit.orchestration.scheduler._pick_resume (локальная копия ладдера — контракт другой: нет params/interval_days, курсор числовой, а не bucket-set, поэтому общий _pick_resume не переиспользован). - Резюмится только если прошлый прогон НЕ завершился штатно: status in ('running', 'zombie'), либо 'done' с counters.interrupted=1 (SIGTERM-drain — см. ветку в import_rosreestr_dkp, которая раньше называла последующий full rescan штатным поведением). - Потолок возраста чекпоинта — 24ч (_DKP_CHECKPOINT_STALE_HOURS): старше — курсор считается негодным, прогон стартует с last_id=0, причина 'checkpoint_stale'. - Вердикт пишется в scrape_runs.counters через update_heartbeat, а тот counters МЕРЖИТ (`counters || :counters`) — иначе первый же per-batch heartbeat после старта стёр бы resume-вердикт. На момент #3168 мерж был только у kit-копии, поэтому запись шла именно kit-именем; с #3390 копия одна и семантика мержа — единственная. Обратимость (см. PR summary): временный откат last_id на литерал 0 красит test_resume_continues_from_saved_last_id (last_id == 123456 не совпадает с 0). """ from __future__ import annotations import inspect import json import os from typing import Any from unittest.mock import MagicMock os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db") from app.services import scheduler as sched class _FakeRow: def __init__(self, **kw: Any) -> None: self.__dict__.update(kw) class _FakeDb: """Двойник сессии: SELECT отдаёт фиксированную строку-кандидата, UPDATE копится. Тот же паттерн, что test_930_scheduler_resume_checkpoint.py::_FakeDb — различает вызовы по наличию ключа 'counters' в params (heartbeat-запись), а не по SQL-тексту. """ def __init__(self, row: Any) -> None: self.row = row self.written: list[dict[str, Any]] = [] def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any: if params is not None and "counters" in params: self.written.append(params) return MagicMock() return MagicMock(fetchone=lambda: self.row) def commit(self) -> None: pass def _candidate(**over: Any) -> _FakeRow: """Строка-кандидат из _DKP_RESUME_CANDIDATE_SQL: прогон 777, убитый на 20-м батче.""" base: dict[str, Any] = { "prev_id": 777, "prev_status": "zombie", "prev_counters": {"rows_fetched": 40000, "last_id": 123456, "batches_done": 20}, "age_h": 2.0, } base.update(over) return _FakeRow(**base) # ── 1. Оборванный прогон продолжает с сохранённого last_id, а не с нуля ───────────── def test_resume_continues_from_saved_last_id() -> None: """status='zombie' (реап после обрыва) + свежий чекпоинт → last_id прошлого прогона. Красный до фикса: import_rosreestr_dkp инициализировал `last_id: int = 0` литералом — эквивалент этого пути (см. docstring модуля § Обратимость и PR summary с реальным выводом pytest до/после отката). """ db = _FakeDb(_candidate(prev_status="zombie", age_h=2.0)) last_id, verdict = sched._resume_dkp_cursor(db, run_id=999) assert last_id == 123456 assert verdict["resume_from"] == 777 assert verdict["resume_reason"] == "ok" assert verdict["last_id"] == 123456 def test_resume_continues_after_sigterm_drain_despite_done_status() -> None: """SIGTERM-drain коммитит status='done' (не 'zombie') — резюмится по counters.interrupted.""" db = _FakeDb( _candidate( prev_status="done", prev_counters={"last_id": 999000, "interrupted": 1}, age_h=0.5, ) ) last_id, verdict = sched._resume_dkp_cursor(db, run_id=1000) assert last_id == 999000 assert verdict["resume_reason"] == "ok" def test_clean_done_run_is_not_resumed() -> None: """Полный проход (status='done', БЕЗ interrupted) — резюмить нечего, старт с 0. Иначе следующий прогон перестал бы пере-сканировать источник целиком и ON CONFLICT DO UPDATE перестал бы ловить правки уже импортированных сделок. """ db = _FakeDb(_candidate(prev_status="done", prev_counters={"last_id": 5000}, age_h=1.0)) last_id, verdict = sched._resume_dkp_cursor(db, run_id=1001) assert last_id == 0 assert verdict["resume_reason"] == "status_done" def test_no_previous_run_starts_from_zero() -> None: """Первый прогон источника вообще — резюмить нечего, явная причина 'no_prev_run'.""" db = _FakeDb(None) last_id, verdict = sched._resume_dkp_cursor(db, run_id=1) assert last_id == 0 assert verdict["resume_reason"] == "no_prev_run" assert verdict["resume_from"] is None # ── 1b. Курсор — ПЕР РЕГИОН (#3051 п.3): source изолирует чекпоинты 66/77 ─────────── class _SourceCapturingDb: """Двойник: запоминает `source`-bind ПОСЛЕДНЕГО SELECT-запроса (_DKP_RESUME_CANDIDATE_SQL). Отличает SELECT (несёт 'source' в params) от heartbeat-UPDATE (несёт 'counters') тем же приёмом, что _FakeDb выше. """ def __init__(self, row: Any) -> None: self.row = row self.select_params: dict[str, Any] | None = None def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any: if params is not None and "source" in params: self.select_params = params return MagicMock(fetchone=lambda: self.row) return MagicMock() def commit(self) -> None: pass def test_resume_uses_region_source() -> None: """source= (#3051 п.3) — курсор региона 77 изолирован от курсора региона 66: явный source попадает bind-параметром в SELECT-кандидата; вызов без аргумента сохраняет прежнее имя _DKP_SOURCE (регион 66, обратная совместимость). Сам call site (import_rosreestr_dkp) обязан передавать именно региональный source, не полагаться на дефолт молча. """ db_77 = _SourceCapturingDb(None) sched._resume_dkp_cursor(db_77, run_id=1, source="rosreestr_dkp_import_77") assert db_77.select_params is not None assert db_77.select_params["source"] == "rosreestr_dkp_import_77" db_default = _SourceCapturingDb(None) sched._resume_dkp_cursor(db_default, run_id=2) assert db_default.select_params is not None assert db_default.select_params["source"] == sched._DKP_SOURCE import_src = inspect.getsource(sched.import_rosreestr_dkp) assert "source=dkp_source" in import_src, "call site не передаёт региональный source" # ── 2. Курсор старше потолка не подхватывается — прогон стартует заново ───────────── def test_stale_checkpoint_is_rejected() -> None: """Чекпоинту больше суток (issue: «курсор недельной давности пропустит правки апстрима ниже last_id») → прогон стартует с 0, причина видна в вердикте.""" db = _FakeDb(_candidate(prev_status="running", prev_counters={"last_id": 42}, age_h=24.01)) last_id, verdict = sched._resume_dkp_cursor(db, run_id=1002) assert last_id == 0 assert verdict["resume_reason"] == "checkpoint_stale" def test_checkpoint_just_under_ceiling_is_still_accepted() -> None: """Граница потолка (23.99ч) — ещё подхватывается, не off-by-one.""" db = _FakeDb(_candidate(prev_status="running", prev_counters={"last_id": 42}, age_h=23.99)) last_id, verdict = sched._resume_dkp_cursor(db, run_id=1003) assert last_id == 42 assert verdict["resume_reason"] == "ok" # ── 3. Запись курсора не затирает посторонние ключи в counters (мерж, не замена) ──── # # Текстового гейта на имя `kit_runs.update_heartbeat` здесь больше нет (#3390): пока копий # было две, имя выбирало семантику, и проверять его по тексту имело смысл. Теперь # `app.services.scrape_runs` — алиас kit'а, оба имени дают ОДИН объект, и гейт краснел бы # на переименовании, ничего при этом не защищая. Мерж проверяется по значению — тестом # ниже и test_3390_single_runs_module.py (оба пути импорта, heartbeat + финализатор). def test_update_heartbeat_merges_into_existing_counters() -> None: """Сама примитива слияния: второй write добавляет rows_fetched/last_id, НЕ стирая resume_from/resume_reason, записанные первым write'ом. SQL-текст реального update_heartbeat (scraper_kit.orchestration.runs) содержит `||` — проверяем это как контракт, а сам мерж эмулируем на стороне двойника (нет живого Postgres в юнит-тесте). """ from scraper_kit.orchestration import runs as kit_runs state: dict[str, Any] = {} class _MergeDb: def execute(self, stmt: Any, params: dict[str, Any]) -> Any: sql = str(stmt) assert "||" in sql, "update_heartbeat должен МЕРЖИТЬ (counters || :counters)" state.update(json.loads(params["counters"])) return MagicMock() def commit(self) -> None: pass db = _MergeDb() kit_runs.update_heartbeat(db, 1, {"resume_from": 777, "resume_reason": "ok"}) kit_runs.update_heartbeat(db, 1, {"rows_fetched": 2000, "last_id": 123456}) assert state["resume_from"] == 777 # не стёрлось вторым write'ом assert state["resume_reason"] == "ok" assert state["rows_fetched"] == 2000 assert state["last_id"] == 123456