Compare commits
No commits in common. "a326bc9b69373b0cb526c3681f7442c7c3c13bda" and "bdb9b64b032724b99e828afd76e9a75b739af031" have entirely different histories.
a326bc9b69
...
bdb9b64b03
2 changed files with 9 additions and 316 deletions
|
|
@ -28,12 +28,6 @@ from __future__ import annotations
|
||||||
import logging
|
import logging
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
# kit_runs.update_heartbeat (в отличие от локального runs_mod.update_heartbeat) мержит
|
|
||||||
# counters (`counters || :counters`) вместо замены — нужен для чекпоинта курсора
|
|
||||||
# import_rosreestr_dkp (issue #3168), чтобы resume-вердикт, записанный на старте, не
|
|
||||||
# затирался последующими per-batch heartbeat'ами того же прогона.
|
|
||||||
from scraper_kit.orchestration import runs as kit_runs
|
|
||||||
|
|
||||||
# compute_next_run_at жил здесь ВТОРОЙ, побайтово одинаковой копией kit-версии (#2674).
|
# compute_next_run_at жил здесь ВТОРОЙ, побайтово одинаковой копией kit-версии (#2674).
|
||||||
# Обе копии одинаково умели interval_days — но такт доезжал до next_run_at только через
|
# Обе копии одинаково умели interval_days — но такт доезжал до next_run_at только через
|
||||||
# kit (_claim_run/_defer_next_run_at читают default_params["interval_days"]); admin.py
|
# kit (_claim_run/_defer_next_run_at читают default_params["interval_days"]); admin.py
|
||||||
|
|
@ -165,85 +159,6 @@ async def _execute_cian_backfill(
|
||||||
raise
|
raise
|
||||||
|
|
||||||
|
|
||||||
_DKP_SOURCE = "rosreestr_dkp_import"
|
|
||||||
|
|
||||||
# Потолок возраста чекпоинта: старше — last_id прошлого прогона не подхватываем, прогон
|
|
||||||
# стартует с id=0 (issue #3168). У предиката `id > last_id` нет протухания в смысле
|
|
||||||
# свипов (он остаётся корректным сколь угодно долго), но апстрим
|
|
||||||
# (gendesign_rosreestr_deals) наполняется НЕЗАВИСИМЫМ ETL, который может дописать более
|
|
||||||
# старую сделку под новым id уже ПОСЛЕ того, как наш курсор её обогнал — курсор недельной
|
|
||||||
# давности унёс бы с собой всё, что апстрим добавил за эту неделю ниже last_id. Порог
|
|
||||||
# того же порядка, что суточная сетка резюма в scraper_kit (_resume_decision), и с
|
|
||||||
# запасом перекрывает самый долгий соседний backfill в семье (avito_detail_backfill,
|
|
||||||
# 150 мин максимума heartbeat-gap за 30 суток — issue #3168).
|
|
||||||
_DKP_CHECKPOINT_STALE_HOURS = 24.0
|
|
||||||
|
|
||||||
_DKP_RESUME_CANDIDATE_SQL = text("""
|
|
||||||
SELECT id AS prev_id,
|
|
||||||
status AS prev_status,
|
|
||||||
counters AS prev_counters,
|
|
||||||
EXTRACT(EPOCH FROM (clock_timestamp() - heartbeat_at)) / 3600.0 AS age_h
|
|
||||||
FROM scrape_runs
|
|
||||||
WHERE source = CAST(:source AS text)
|
|
||||||
AND id <> CAST(:rid AS bigint)
|
|
||||||
AND status <> 'skipped'
|
|
||||||
ORDER BY started_at DESC
|
|
||||||
LIMIT 1
|
|
||||||
""")
|
|
||||||
|
|
||||||
|
|
||||||
def _resume_dkp_cursor(db: Session, run_id: int) -> tuple[int, dict[str, Any]]:
|
|
||||||
"""Продолжить last_id прошлого прогона или начать с 0 — решение + explain (issue #3168).
|
|
||||||
|
|
||||||
last_id раньше жил только в памяти процесса (init на 0 при каждом запуске): heartbeat
|
|
||||||
писал его в counters КАЖДЫЙ батч (checkpoint), но на старте никто это не читал
|
|
||||||
обратно. Обрыв (деплой/OOM/рестарт хоста) откатывал прогресс на 0 и заставлял
|
|
||||||
пере-сканировать источник с начала.
|
|
||||||
|
|
||||||
Кандидат — ПОСЛЕДНИЙ прогон этого source (тот же принцип, что и
|
|
||||||
scraper_kit.orchestration.scheduler._pick_resume, локальная копия ладдера — контракт
|
|
||||||
другой: нет params/interval_days, курсор числовой, а не bucket-set):
|
|
||||||
- 'running' / 'zombie' — прогон, которого не завершили штатно.
|
|
||||||
- 'done' с counters.interrupted=1 — SIGTERM-drain (см. комментарий у mark_done
|
|
||||||
ниже по коду): статус 'done', но это НЕ полный проход, прогресс оборван.
|
|
||||||
- Чистое 'done' без этого флага — полный проход завершился сам, резюмить нечего:
|
|
||||||
следующий прогон обязан пере-сканировать с 0, иначе ON CONFLICT DO UPDATE
|
|
||||||
перестанет ловить правки уже импортированных сделок (см. докстринг функции ниже).
|
|
||||||
|
|
||||||
Возвращает (last_id, verdict) — verdict пишется в scrape_runs.counters вызывающим
|
|
||||||
кодом (kit_runs.update_heartbeat — merge, не замена), чтобы решение было видно в
|
|
||||||
scrape_runs, а не только в логе.
|
|
||||||
"""
|
|
||||||
row = db.execute(_DKP_RESUME_CANDIDATE_SQL, {"source": _DKP_SOURCE, "rid": run_id}).fetchone()
|
|
||||||
verdict: dict[str, Any] = {"resume_from": None}
|
|
||||||
|
|
||||||
if row is None:
|
|
||||||
verdict["resume_reason"] = "no_prev_run"
|
|
||||||
return 0, verdict
|
|
||||||
|
|
||||||
verdict["resume_candidate"] = int(row.prev_id)
|
|
||||||
prev_counters = row.prev_counters if isinstance(row.prev_counters, dict) else {}
|
|
||||||
prev_last_id = prev_counters.get("last_id")
|
|
||||||
interrupted = bool(prev_counters.get("interrupted"))
|
|
||||||
resumable_status = row.prev_status in ("running", "zombie") or (
|
|
||||||
row.prev_status == "done" and interrupted
|
|
||||||
)
|
|
||||||
|
|
||||||
if not resumable_status:
|
|
||||||
verdict["resume_reason"] = f"status_{row.prev_status}"
|
|
||||||
elif not isinstance(prev_last_id, int) or prev_last_id <= 0:
|
|
||||||
verdict["resume_reason"] = "no_checkpoint"
|
|
||||||
elif row.age_h is None or float(row.age_h) > _DKP_CHECKPOINT_STALE_HOURS:
|
|
||||||
verdict["resume_reason"] = "checkpoint_stale"
|
|
||||||
else:
|
|
||||||
verdict["resume_from"] = int(row.prev_id)
|
|
||||||
verdict["resume_reason"] = "ok"
|
|
||||||
verdict["last_id"] = prev_last_id
|
|
||||||
return int(prev_last_id), verdict
|
|
||||||
|
|
||||||
return 0, verdict
|
|
||||||
|
|
||||||
|
|
||||||
def import_rosreestr_dkp(
|
def import_rosreestr_dkp(
|
||||||
db: Session,
|
db: Session,
|
||||||
run_id: int,
|
run_id: int,
|
||||||
|
|
@ -280,10 +195,7 @@ def import_rosreestr_dkp(
|
||||||
Rooms: выводятся из площади (Росреестр не отдаёт кол-во комнат).
|
Rooms: выводятся из площади (Росреестр не отдаёт кол-во комнат).
|
||||||
|
|
||||||
Batch-процессинг: читаем из FDW батчами по batch_size через cursor-based пагинацию
|
Batch-процессинг: читаем из FDW батчами по batch_size через cursor-based пагинацию
|
||||||
(WHERE id > last_id ORDER BY id). Heartbeat обновляется каждый батч (= checkpoint),
|
(WHERE id > last_id ORDER BY id). Heartbeat обновляется каждый батч (= checkpoint).
|
||||||
мержем (kit_runs.update_heartbeat), а не заменой. На старте _resume_dkp_cursor решает
|
|
||||||
продолжить с last_id прошлого прогона или начать с 0 — чекпоинт переживает рестарт
|
|
||||||
процесса (деплой/OOM/SIGTERM), пока не старше суток (issue #3168).
|
|
||||||
SAVEPOINT per row — один сбойный row не откатывает батч.
|
SAVEPOINT per row — один сбойный row не откатывает батч.
|
||||||
|
|
||||||
Координаты: NULL после импорта — геокодинг остаётся follow-up (geocode-deals).
|
Координаты: NULL после импорта — геокодинг остаётся follow-up (geocode-deals).
|
||||||
|
|
@ -329,14 +241,8 @@ def import_rosreestr_dkp(
|
||||||
)
|
)
|
||||||
db.rollback()
|
db.rollback()
|
||||||
|
|
||||||
last_id, resume_verdict = _resume_dkp_cursor(db, run_id)
|
last_id: int = 0
|
||||||
total_batches = 0
|
total_batches = 0
|
||||||
kit_runs.update_heartbeat(db, run_id, resume_verdict)
|
|
||||||
logger.info(
|
|
||||||
"rosreestr_dkp_import run_id=%d: resume decision — %s",
|
|
||||||
run_id,
|
|
||||||
resume_verdict,
|
|
||||||
)
|
|
||||||
|
|
||||||
try:
|
try:
|
||||||
while True:
|
while True:
|
||||||
|
|
@ -349,14 +255,11 @@ def import_rosreestr_dkp(
|
||||||
return
|
return
|
||||||
# #1182 Phase 2: кооперативный SIGTERM-drain (деплой recreate scraper).
|
# #1182 Phase 2: кооперативный SIGTERM-drain (деплой recreate scraper).
|
||||||
# Это НЕ user-cancel → mark_done (partial), не mark_cancelled. Курсор
|
# Это НЕ user-cancel → mark_done (partial), не mark_cancelled. Курсор
|
||||||
# уже зафиксирован heartbeat'ом каждый батч; следующий run подхватит
|
# уже зафиксирован update_heartbeat'ом каждый батч; следующий run
|
||||||
# last_id через _resume_dkp_cursor (issue #3168), если чекпоинт не старше
|
# пере-сканирует с id=0 — INSERT ... ON CONFLICT DO UPDATE идемпотентен
|
||||||
# суток — раньше здесь безусловно пере-сканировали с id=0 при каждом
|
# для неизменных строк (WHERE ... IS DISTINCT FROM пропускает совпадающие
|
||||||
# обрыве. counters['interrupted']=1 помечает это 'done' как НЕ полный
|
# факты), поэтому повторный проход не плодит дубликаты и не трогает данные.
|
||||||
# проход, чтобы _resume_dkp_cursor не спутал его со штатным завершением
|
|
||||||
# (которое обязано пере-сканировать с 0 ради ON CONFLICT DO UPDATE правок).
|
|
||||||
# mark_done выводит run из 'running' → reap_zombies его не тронет.
|
# mark_done выводит run из 'running' → reap_zombies его не тронет.
|
||||||
counters["interrupted"] = 1
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"rosreestr_dkp_import run_id=%d: SIGTERM-drain — committing partial "
|
"rosreestr_dkp_import run_id=%d: SIGTERM-drain — committing partial "
|
||||||
"(last_id=%d, batches=%d) and exiting",
|
"(last_id=%d, batches=%d) and exiting",
|
||||||
|
|
@ -364,7 +267,7 @@ def import_rosreestr_dkp(
|
||||||
last_id,
|
last_id,
|
||||||
total_batches,
|
total_batches,
|
||||||
)
|
)
|
||||||
kit_runs.update_heartbeat(db, run_id, counters)
|
runs_mod.update_heartbeat(db, run_id, counters)
|
||||||
runs_mod.mark_done(db, run_id, counters)
|
runs_mod.mark_done(db, run_id, counters)
|
||||||
return
|
return
|
||||||
|
|
||||||
|
|
@ -538,10 +441,8 @@ def import_rosreestr_dkp(
|
||||||
counters["batches_done"] = total_batches
|
counters["batches_done"] = total_batches
|
||||||
counters["last_id"] = last_id # type: ignore[assignment]
|
counters["last_id"] = last_id # type: ignore[assignment]
|
||||||
|
|
||||||
# Heartbeat = checkpoint: allows zombie detection + resume visibility.
|
# Heartbeat = checkpoint: allows zombie detection + resume visibility
|
||||||
# kit_runs (merge, не замена) — иначе этот REPLACE стёр бы resume_verdict,
|
runs_mod.update_heartbeat(db, run_id, counters)
|
||||||
# записанный _resume_dkp_cursor'ом перед циклом (issue #3168).
|
|
||||||
kit_runs.update_heartbeat(db, run_id, counters)
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"rosreestr_dkp_import run_id=%d: batch=%d fetched=%d "
|
"rosreestr_dkp_import run_id=%d: batch=%d fetched=%d "
|
||||||
"inserted=%d updated=%d skipped=%d errored=%d last_id=%d",
|
"inserted=%d updated=%d skipped=%d errored=%d last_id=%d",
|
||||||
|
|
|
||||||
|
|
@ -1,208 +0,0 @@
|
||||||
"""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 через kit_runs.update_heartbeat
|
|
||||||
(`counters || :counters` — merge), а не локальный runs_mod.update_heartbeat
|
|
||||||
(`CAST(:counters AS jsonb)` — полная замена): иначе первый же per-batch heartbeat
|
|
||||||
после старта стирает resume-вердикт.
|
|
||||||
|
|
||||||
Обратимость (см. 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
|
|
||||||
|
|
||||||
|
|
||||||
# ── 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 (мерж, не замена) ────
|
|
||||||
|
|
||||||
|
|
||||||
def test_cursor_write_uses_merge_not_replace_heartbeat() -> None:
|
|
||||||
"""import_rosreestr_dkp обязан писать чекпоинт через kit_runs.update_heartbeat
|
|
||||||
(merge: `counters || :counters`), а не локальный runs_mod.update_heartbeat (замена:
|
|
||||||
`CAST(:counters AS jsonb)`) — иначе resume-вердикт, записанный ДО цикла, стирается
|
|
||||||
первым же per-batch heartbeat'ом того же прогона.
|
|
||||||
"""
|
|
||||||
src = inspect.getsource(sched.import_rosreestr_dkp)
|
|
||||||
assert "kit_runs.update_heartbeat" in src
|
|
||||||
assert "runs_mod.update_heartbeat" not in src
|
|
||||||
|
|
||||||
|
|
||||||
def test_kit_runs_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
|
|
||||||
Loading…
Add table
Reference in a new issue