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
237 lines
12 KiB
Python
237 lines
12 KiB
Python
"""#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 + порог сравнения
|