"""#2702: отметки времени прогона не охватывают его работу. Что было. Финализаторы писали `finished_at`/`heartbeat_at` через `now()`, а `now()` в PostgreSQL — синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции. Выполняются финализаторы той же сессией, что и работа задачи, поэтому если рабочая транзакция всё это время оставалась открытой, их UPDATE попадал ВНУТРЬ неё и получал время НАЧАЛА работы. Прод-замер 2026-08-06 (487 прогонов с finished_at и counters.duration_sec): * 153 — заявленная длительность больше окна finished_at − started_at в >1.5 раза; * 133 — окно меньше секунды при работе дольше 10 с, причём 126 из них укладываются в 9-64 мс: столько проходит от коммита claim'а до первого запроса рабочей транзакции — подпись механизма, а не разброс. Крайний случай, воспроизведённый ниже дословно: прогон 346 (cian_history_backfill) — 18 230 с работы, окно 32 мс. Почему дефект был не сплошной: он зависел от того, коммитила ли задача перед финалом. `cadastral_geo_match` / `house_imv_backfill` / `avito_detail_backfill` коммитят поштучно — у них окно совпадало с работой; `yandex_address_backfill` (45 из 50 прогонов), `newbuilding_enrich`, `cian_history_backfill` — нет. Фальсификация: на старом коде (`now()`) тесты 1 и 3 дают другой ответ. """ from __future__ import annotations import inspect import os import re from typing import Any import pytest os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db") from scraper_kit.orchestration import runs as kit_runs from scraper_kit.orchestration import scheduler as kit_scheduler from app.services import scrape_runs as app_runs _MODULES = {"kit": kit_runs, "app": app_runs} # Колонки, которые обязаны нести НАСТОЯЩЕЕ время, а не время старта транзакции. _TS_COLS = ("started_at", "finished_at", "heartbeat_at") def _split_top(items: str) -> list[str]: """Разбить список SQL-элементов по запятым ВЕРХНЕГО уровня (CAST(:x AS t) — один).""" out: list[str] = [] depth = 0 cur = "" for ch in items: if ch == "," and depth == 0: out.append(cur.strip()) cur = "" continue depth += (ch == "(") - (ch == ")") cur += ch out.append(cur.strip()) return out class _Row: """Строка ответа: id для create_run, source для alert-хука.""" id = 1 source = "src" status = "done" class _FakeResult: def fetchone(self) -> _Row: return _Row() def first(self) -> _Row: return _Row() def fetchall(self) -> list[_Row]: return [] class _FakePg: """Мини-модель PostgreSQL на две функции времени и ленивую транзакцию. `now()` == `transaction_timestamp()` — замерзает на старте транзакции; `clock_timestamp()` — настоящие часы. Транзакция открывается лениво на первом execute (autobegin SQLAlchemy) и закрывается commit/rollback. Больше модель ничего не умеет — этого достаточно, чтобы отличить одно от другого. """ def __init__(self) -> None: self.wall: float = 0.0 # «стенные часы» теста self.tx_start: float | None = None self.row: dict[str, float] = {} # что осело в scrape_runs self.calls: list[str] = [] def _stamp(self, col: str, func: str) -> None: assert self.tx_start is not None self.row[col] = self.tx_start if func.lower() == "now" else self.wall def execute(self, stmt: Any, params: Any = None) -> _FakeResult: if self.tx_start is None: self.tx_start = self.wall self.calls.append("execute") sql = str(stmt) for col in _TS_COLS: # UPDATE ... SET = () m = re.search(rf"\b{col}\s*=\s*(now|clock_timestamp)\s*\(\s*\)", sql, re.I) if m is not None: self._stamp(col, m.group(1)) ins = re.search( r"INSERT INTO scrape_runs\s*\((.*?)\).*?VALUES\s*\((.*)\)", sql, re.S | re.I ) if ins is not None: # INSERT — отметки времени позиционные, в VALUES for col, val in zip(_split_top(ins.group(1)), _split_top(ins.group(2)), strict=False): f = re.fullmatch(r"(now|clock_timestamp)\s*\(\s*\)", val, re.I) if col in _TS_COLS and f is not None: self._stamp(col, f.group(1)) return _FakeResult() def commit(self) -> None: self.calls.append("commit") self.tx_start = None def rollback(self) -> None: self.calls.append("rollback") self.tx_start = None # ── 1. Финал прогона датируется концом работы, а не стартом транзакции ──────── @pytest.mark.parametrize("name", list(_MODULES)) def test_finished_at_covers_the_work_not_the_transaction_start(name: str) -> None: """Прогон 346 дословно: 18 230 с работы в одной незакоммиченной транзакции. На старом коде finished_at = 0.032 (старт рабочей транзакции) → окно 32 мс при пяти часах работы. Это и есть та строка, ради которой заведена задача. """ mod = _MODULES[name] db = _FakePg() db.row["started_at"] = 0.0 # create_run уже закоммитил claim db.wall = 0.020 mod.update_heartbeat(db, 1, {"listings_processed": 0}) # heartbeat + commit db.wall = 0.032 db.execute("SELECT id FROM listings WHERE history IS NULL") # рабочая транзакция db.wall = 18230.0 # пять часов работы, ни одного коммита mod.mark_done(db, 1, {"listings_processed": 1200, "duration_sec": 18230}) assert db.row["finished_at"] == pytest.approx(18230.0) assert db.row["heartbeat_at"] == pytest.approx(18230.0) window = db.row["finished_at"] - db.row["started_at"] assert window == pytest.approx(18230.0), "окно прогона обязано охватывать его работу" @pytest.mark.parametrize("name", list(_MODULES)) @pytest.mark.parametrize("finalizer", ["mark_failed", "mark_banned"]) def test_failed_and_banned_finals_are_wall_clock_too(name: str, finalizer: str) -> None: """Тот же инвариант для неуспешных финалов. Их спасал defensive-rollback в начале (он закрывал рабочую транзакцию), но полагаться на побочный эффект чужой защиты нельзя — проверяем явно. """ mod = _MODULES[name] db = _FakePg() db.wall = 0.019 db.execute("SELECT 1") # рабочая транзакция открыта db.wall = 1460.0 getattr(mod, finalizer)(db, 1, "boom", {"checked": 0, "duration_sec": 1460}) assert db.row["finished_at"] == pytest.approx(1460.0) # ── 2. started_at переживает откат рабочей транзакции ──────────────────────── @pytest.mark.parametrize("name", list(_MODULES)) def test_started_at_survives_rolled_back_work_transaction(name: str) -> None: """Требование #2702 п.1: отметка старта живёт в СВОЕЙ закоммиченной транзакции. create_run коммитит INSERT до возврата run_id, а ни один финализатор прогона started_at не переписывает — поэтому откат рабочей транзакции его не достаёт. Финал при этом обязан быть датирован концом работы (это и падает на старом коде). """ mod = _MODULES[name] db = _FakePg() run_id = mod.create_run(db, source="cian_history_backfill", params={}) assert run_id == 1 assert db.calls == ["execute", "commit"], "INSERT прогона обязан коммититься сразу" db.wall = 0.030 db.execute("UPDATE listings SET address = 'x'") # рабочая транзакция db.wall = 100.0 db.rollback() # работа упала и откатилась db.wall = 100.5 db.execute("SELECT count(*) FROM listings") # новая рабочая транзакция db.wall = 1460.0 mod.mark_done(db, 1, {"checked": 0, "duration_sec": 1460}) assert db.row["started_at"] == pytest.approx(0.0), "started_at не должен сдвигаться" assert db.row["finished_at"] == pytest.approx(1460.0) # ── 3. Инвариант источника: никакая отметка времени не пишется now() ───────── @pytest.mark.parametrize("name", list(_MODULES)) def test_no_run_timestamp_is_written_with_now(name: str) -> None: """`now()` в этих модулях не имеет корректного применения — его быть не должно. Проверяем весь исходник, а не отдельные запросы: INSERT в create_run пишет started_at/heartbeat_at позиционно (в VALUES), и построчная проверка его бы пропустила — ровно так дефект и дожил до 3 300 прогонов. """ src = inspect.getsource(_MODULES[name]) # Регистрозависимо: SQL в этих модулях пишется в верхнем регистре, а строчное # `now()` встречается в объяснительной прозе docstring'ов — ловим SQL, не текст. assert re.search(r"\bNOW\s*\(\s*\)", src) is None assert "clock_timestamp()" in src def test_zombie_criterion_compares_real_clocks() -> None: """#2702 п.2: поиск зависших сравнивает записанный heartbeat со «сейчас». Обе стороны сравнения обязаны быть настоящим временем: на проде у всех 6 прогонов cian_history_backfill, помеченных 'zombie', записанный heartbeat так и остался на отметке старта (max advance 0.0 с) — критерий решал по замороженной отметке, хотя нормальный прогон этого источника длится до 5.06 ч. """ src = inspect.getsource(kit_scheduler.reap_zombies) stmt = re.search(r"UPDATE scrape_runs.*?RETURNING id", src, re.S) assert stmt is not None assert re.search(r"\bNOW\s*\(\s*\)", stmt.group(0)) is None assert stmt.group(0).count("clock_timestamp()") == 2 # finished_at + порог сравнения