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
292 lines
14 KiB
Python
292 lines
14 KiB
Python
"""#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"
|