"""#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"