gendesign/tradein-mvp/backend/tests/test_2702_run_timestamps.py
bot-backend 663a831775
Some checks failed
Deploy Trade-In / test (push) Blocked by required conditions
Deploy Trade-In / build-backend (push) Blocked by required conditions
Deploy Trade-In / build-frontend (push) Blocked by required conditions
Deploy Trade-In / build-browser (push) Blocked by required conditions
Deploy Trade-In / deploy (push) Blocked by required conditions
Deploy Trade-In / changes (push) Has been cancelled
fix(tradein/scraper): отметки времени прогона перестают замерзать в его же транзакции (#2702) (#2718)
2026-08-06 09:55:50 +00:00

237 lines
12 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""#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 <col> = <func>()
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 + порог сравнения