gendesign/tradein-mvp/backend/tests/test_2731_observation_timestamps.py
bot-backend 91423e0b53
Some checks failed
Deploy Trade-In / changes (push) Successful in 24s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m41s
Deploy Trade-In / build-backend (push) Successful in 1m56s
Deploy Trade-In / deploy (push) Has been cancelled
fix(tradein/scrapers): метка наблюдения — время строки, а не старта транзакции (#2731) (#2742)
2026-08-06 16:19:56 +00:00

292 lines
14 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.

"""#2731: метка наблюдения — время строки, а не старта транзакции сбора.
Что было. `save_listings` и `upsert_listing_snapshot` писали содержательные колонки
времени через `NOW()`, а `NOW()` в PostgreSQL — синоним `transaction_timestamp()`:
он замерзает на СТАРТЕ транзакции. Коммит у писателя объявлений ОДИН, в конце всего
batch'а, поэтому все строки вызова получали одну и ту же метку — сколько бы он ни работал.
Прод-замер 2026-08-06 (listings_snapshots, строк / различных меток на прогон):
run 3303 — 219 / 1, при семи минутах работы (11:59:59 → 12:07:19);
run 3299 — 235 / 1; run 3293 — 297 / 1;
run 3229 — 2643 / 55 — метка на ВЫЗОВ save_listings, а не на строку.
`listings.scraped_at` и `last_seen_at` за те же сутки схлопнуты так же и при этом равны
друг другу у 100% строк (по часам: eq == rows во всех корзинах).
Цена дефекта — НЕ фильтры свежести (смещение равно длительности прогона: минуты, худший
случай 17.3 ч, против окна 14 суток), а разрешение во времени: темп сбора и порядок строк
внутри прогона по данным не восстановить (отсюда же невозможность бэкфилла run_id, #2701).
Почему `statement_timestamp()`, а не `clock_timestamp()` как в #2702/#2718: здесь в одном
statement'е пишутся ДВЕ колонки, обязанные совпадать (после #2206 scraped_at = last_seen_at,
на этом равенстве стоит предикат миграции 161). На проде проверено:
`clock_timestamp() = clock_timestamp()` → false, `statement_timestamp() = ...` → true.
Модель БД ниже воспроизводит ровно это различие, поэтому «починка» через clock_timestamp()
провалит тест на согласованность.
Фальсификация: на старом коде (`NOW()`) тесты 1 и 2 дают одну метку на все строки.
"""
from __future__ import annotations
import inspect
import os
import re
from contextlib import contextmanager
from pathlib import Path
from typing import Any
from unittest.mock import MagicMock
import pytest
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
from scraper_kit import base as kit_base
from scraper_kit import snapshot_writer as kit_snapshot
from scraper_kit.base import ScrapedLot, save_listings
# Сколько «работает» писатель на одну строку. Пять лотов → 2 с работы одной транзакции:
# требование задачи — писатель, отработавший ДОЛЬШЕ СЕКУНДЫ, оставляет различные метки.
_STEP_SEC = 0.4
# Содержательные колонки времени, за которыми следим, по таблицам-писателям.
_WATCHED: dict[str, tuple[str, ...]] = {
"listings": ("scraped_at", "last_seen_at"),
"listings_snapshots": ("observed_at",),
}
def _strip_sql_comments(sql: str) -> str:
"""Убрать `--`-комментарии: в них есть и скобки, и слово NOW() в объяснительной прозе."""
return re.sub(r"--[^\n]*", "", sql)
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
def _balanced(sql: str, open_idx: int) -> tuple[str, int]:
"""Содержимое скобки, открытой на open_idx, и индекс её закрывающей пары."""
depth = 0
for i in range(open_idx, len(sql)):
depth += (sql[i] == "(") - (sql[i] == ")")
if depth == 0:
return sql[open_idx + 1 : i], i
raise AssertionError("несбалансированные скобки в SQL")
class _Row:
"""Строка ответа: id/inserted для INSERT INTO listings ... RETURNING."""
id = 42
inserted = False
card_hash = None
class _FakeResult:
def __init__(self, row: _Row | None) -> None:
self._row = row
def fetchone(self) -> _Row | None:
return self._row
def first(self) -> _Row | None:
return self._row
def scalar_one_or_none(self) -> None:
return None
class _FakePg:
"""Мини-модель PostgreSQL на три функции времени и ленивую транзакцию.
`now()` == `transaction_timestamp()` — замерзает на старте транзакции (autobegin на
первом execute, сброс на commit/rollback);
`statement_timestamp()` — момент начала ТЕКУЩЕГО statement'а: двигается от запроса к
запросу и СТАБИЛЕН внутри запроса;
`clock_timestamp()` — настоящие часы, вычисляются на КАЖДЫЙ вызов отдельно, поэтому
два вызова в одном statement'е дают разные значения (проверено на проде).
Каждый execute «работает» _STEP_SEC — так модель отличает построчную метку от общей.
"""
def __init__(self) -> None:
self.wall: float = 0.0
self.tx_start: float | None = None
self.stmt_start: float = 0.0
self._clock_calls: int = 0
# Что осело в таблицах: список строк {колонка: значение}.
self.rows: dict[str, list[dict[str, float]]] = {t: [] for t in _WATCHED}
# ── модель времени ────────────────────────────────────────────────────────
def _value(self, func: str) -> float:
if func == "now":
assert self.tx_start is not None
return self.tx_start
if func == "statement_timestamp":
return self.stmt_start
if func == "clock_timestamp":
self._clock_calls += 1 # каждый вызов — своё значение
return self.stmt_start + self._clock_calls * 1e-6
raise AssertionError(f"неизвестная функция времени: {func}")
# ── разбор писателя ───────────────────────────────────────────────────────
def _record(self, sql: str) -> None:
ins = re.search(r"INSERT INTO (listings|listings_snapshots)\s*\(", sql)
if ins is not None:
table = ins.group(1)
cols, end = _balanced(sql, ins.end() - 1)
vals_kw = re.compile(r"\bVALUES\s*\(").search(sql, end)
assert vals_kw is not None, "INSERT без VALUES"
vals, _ = _balanced(sql, vals_kw.end() - 1)
row: dict[str, float] = {}
for col, val in zip(_split_top(cols), _split_top(vals), strict=False):
f = re.fullmatch(r"(\w+)\s*\(\s*\)", val)
if col in _WATCHED[table] and f is not None:
row[col] = self._value(f.group(1).lower())
if row:
self.rows[table].append(row)
return
upd = re.search(r"UPDATE (listings)\b", sql)
if upd is not None: # reconcile-путь при дрейфе dedup_hash
table = upd.group(1)
row = {}
for col in _WATCHED[table]:
m = re.search(rf"\b{col}\s*=\s*(\w+)\s*\(\s*\)", sql)
if m is not None:
row[col] = self._value(m.group(1).lower())
if row:
self.rows[table].append(row)
# ── интерфейс сессии ──────────────────────────────────────────────────────
def execute(self, stmt: Any, params: Any = None) -> _FakeResult:
if self.tx_start is None:
self.tx_start = self.wall
self.stmt_start = self.wall
self.wall += _STEP_SEC # statement отработал
sql = _strip_sql_comments(str(stmt))
self._record(sql)
if "INSERT INTO listings (" in sql:
return _FakeResult(_Row())
return _FakeResult(None)
def commit(self) -> None:
self.tx_start = None
def rollback(self) -> None:
self.tx_start = None
def begin_nested(self) -> Any:
@contextmanager
def _ctx() -> Any:
yield MagicMock()
return _ctx()
def _matcher() -> MagicMock:
matcher = MagicMock()
matcher.match_or_create_house.return_value = (101, 1.0, "new")
matcher.upsert_listing_source.return_value = None
return matcher
def _lots(n: int = 5) -> list[ScrapedLot]:
return [
ScrapedLot(
source="cian",
source_url=f"https://ekb.cian.ru/sale/flat/{i}/",
source_id=str(i),
price_rub=5_000_000 + i,
)
for i in range(n)
]
def _run_writer() -> _FakePg:
db = _FakePg()
save_listings(db, _lots(), matcher=_matcher(), region_code=66)
return db
# ── 1. Снимки одного batch'а датируются построчно ────────────────────────────
def test_snapshot_marks_differ_per_row() -> None:
"""Прогон 3303 дословно: строки одного вызова — разные моменты, а не один.
На старом коде (`NOW()`) все снимки batch'а несут метку старта транзакции — ровно
219 строк на одну метку, ради чего заведена задача.
"""
db = _run_writer()
marks = [row["observed_at"] for row in db.rows["listings_snapshots"]]
assert len(marks) == 5, "снимок пишется на каждый сохранённый лот"
assert max(marks) - min(marks) > 1.0, "писатель отработал дольше секунды"
assert len(set(marks)) == len(marks), "у строк одного batch'а обязаны быть разные метки"
# ── 2. scraped_at/last_seen_at: построчно, но по-прежнему равны между собой ──
def test_listing_marks_differ_per_row_and_stay_equal_to_each_other() -> None:
"""Две колонки одного statement'а: построчная метка без нового расхождения.
Равенство scraped_at = last_seen_at — не косметика: на нём стоит предикат миграции
161 (`last_seen_at > scraped_at` как признак «видели живым, но не пере-скрейпили»).
`clock_timestamp()` вычисляется на каждый вызов отдельно и развёл бы колонки на
микросекунды — модель БД это воспроизводит, поэтому такая «починка» провалит тест.
"""
db = _run_writer()
rows = db.rows["listings"]
assert len(rows) == 5
for row in rows:
assert row["scraped_at"] == row["last_seen_at"], "колонки statement'а обязаны совпадать"
scraped = [row["scraped_at"] for row in rows]
assert max(scraped) - min(scraped) > 1.0
assert len(set(scraped)) == len(scraped), "метки строк одного batch'а обязаны различаться"
# ── 3. Инвариант источника: содержательная колонка не пишется now() ─────────
@pytest.mark.parametrize("mod", [kit_base, kit_snapshot], ids=["base", "snapshot_writer"])
def test_no_content_timestamp_is_written_with_now(mod: Any) -> None:
"""`<колонка> = NOW()` в этих писателях не имеет корректного применения.
Комментарии вырезаются: `NOW()` там остаётся в объяснительной прозе (и должен —
без неё следующий читатель повторит дефект).
"""
src = _strip_sql_comments(inspect.getsource(mod))
cols = "|".join(sorted({c for t in _WATCHED.values() for c in t}))
assert re.search(rf"\b({cols})\s*=\s*NOW\s*\(\s*\)", src, re.I) is None
assert "statement_timestamp()" in src
# ── 4. Историческая граница зафиксирована в схеме ────────────────────────────
_SQL_DIR = Path(__file__).resolve().parents[1] / "data" / "sql"
_MIGRATION = _SQL_DIR / "232_listings_observation_time_meaning.sql"
def test_migration_232_marks_the_transition_date_and_unrecoverable_history() -> None:
"""Аналитике нужно знать, где сменился смысл колонки и что до него не чинится."""
sql = _MIGRATION.read_text("utf-8")
assert "BEGIN;" in sql and "COMMIT;" in sql
for col in ("listings_snapshots.observed_at", "listings.scraped_at", "listings.last_seen_at"):
assert f"COMMENT ON COLUMN {col} IS" in sql
assert sql.count("#2731 (2026-08-06)") >= 3, "дата перехода — в каждом комментарии"
assert "невосстановим" in sql
assert not re.search(r":\w+::", sql), "psycopg v3: никаких :param::type"