All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m2s
Тесты пишутся ПОСЛЕ живой проверки функционала (правило от 08.09): импорт по 77 — run 6460 (212 937 строк), регион-фильтр коридора — smoke до/после импорта. 14 тестов: SQL-текст через inspect (regex \s+ для многострочных клозов), чистые функции, мок db по образцу 2846, subprocess bash с урезанным PATH (падение от валидации, не от отсутствия docker).
249 lines
13 KiB
Python
249 lines
13 KiB
Python
"""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
|