"""#2993: listing_source_snapshots получает строку только при изменении источника. Проверки идут по ЗНАЧЕНИЯМ в живом Postgres: какие источники получили строку за сегодня, с какими значениями, и какие события записал event-diff поверх них. Без реальной БД файл self-skip'ается, как соседние live-тесты (в CI есть postgres, см. ci-tradein.yml). Вся синтетика живёт в одной транзакции, очистка — rollback. """ from __future__ import annotations import os import uuid from typing import Any os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") import pytest from sqlalchemy import text from app.tasks import listing_source_snapshot as snap_mod def _live_session() -> Any | None: """Тот же контракт, что у test_revisit_floor_lateral_lookup.py.""" try: from sqlalchemy import create_engine from sqlalchemy.orm import sessionmaker dsn = os.environ.get("TEST_DATABASE_URL") or os.environ.get("DATABASE_URL", "") if not dsn or "localhost:5432/test" in dsn: return None engine = create_engine(dsn, future=True) conn = engine.connect() conn.execute(text("SELECT 1")) conn.close() return sessionmaker(bind=engine, future=True)() except Exception: return None pytestmark = pytest.mark.skipif(_live_session() is None, reason="no reachable Postgres test DB") _PRICE = 5_000_000 def _source(db: Any, tag: str, *, price: int, payload: str, seen_days_ago: int) -> int: ext_id = f"zzz2993-{tag}-{uuid.uuid4().hex[:8]}" listing_id = db.execute( text( """ INSERT INTO listings (source, source_url, source_id, dedup_hash, price_rub, is_active, scraped_at, last_seen_at) VALUES ('cian', :url, :ext_id, :ext_id, :price, true, now(), now()) RETURNING id """ ), {"url": f"https://example.test/2993/{ext_id}", "ext_id": ext_id, "price": price}, ).scalar_one() return int( db.execute( text( """ INSERT INTO listing_sources (listing_id, ext_source, ext_id, confidence, matched_method, price_rub, raw_payload, last_seen_at) VALUES (:listing_id, 'cian', :ext_id, 1.0, 'test', :price, CAST(:payload AS jsonb), now() - make_interval(days => :seen)) RETURNING id """ ), { "listing_id": listing_id, "ext_id": ext_id, "price": price, "payload": payload, "seen": seen_days_ago, }, ).scalar_one() ) def _snapshot( db: Any, lsid: int, *, days_ago: int, price: int, payload: str, seen_days_ago: int, is_active: bool = True, ) -> None: # now() постоянен внутри транзакции — last_seen_at снимка совпадает с # listing_sources.last_seen_at до микросекунды, если seen_days_ago тот же. db.execute( text( """ INSERT INTO listing_source_snapshots (listing_source_id, snapshot_date, price_rub, is_active, last_seen_at, payload_hash) VALUES (:lsid, CURRENT_DATE - :days_ago, :price, :is_active, now() - make_interval(days => :seen), md5(CAST(:payload AS jsonb)::text)) """ ), { "lsid": lsid, "days_ago": days_ago, "price": price, "is_active": is_active, "seen": seen_days_ago, "payload": payload, }, ) def _run_writer(db: Any) -> None: db.execute( snap_mod._SNAPSHOT_SQL, {"freshness_days": snap_mod.FRESHNESS_WINDOW_DAYS, "run_id": None}, ) def _today(db: Any, ids: list[int]) -> dict[int, tuple[int, bool]]: rows = db.execute( text( """ SELECT listing_source_id, price_rub, is_active FROM listing_source_snapshots WHERE snapshot_date = CURRENT_DATE AND listing_source_id = ANY(:ids) """ ), {"ids": ids}, ).fetchall() return {int(r[0]): (int(r[1]), bool(r[2])) for r in rows} def _seed(db: Any) -> dict[str, int]: """Семь источников: по одному на каждую ветку условия записи.""" p = '{"v": 1}' ids = { "same": _source(db, "same", price=_PRICE, payload=p, seen_days_ago=1), "same_gap": _source(db, "same_gap", price=_PRICE, payload=p, seen_days_ago=1), "price": _source(db, "price", price=_PRICE, payload=p, seen_days_ago=1), "payload": _source(db, "payload", price=_PRICE, payload=p, seen_days_ago=1), "seen": _source(db, "seen", price=_PRICE, payload=p, seen_days_ago=1), "active": _source(db, "active", price=_PRICE, payload=p, seen_days_ago=10), "new": _source(db, "new", price=_PRICE, payload=p, seen_days_ago=1), } _snapshot(db, ids["same"], days_ago=1, price=_PRICE, payload=p, seen_days_ago=1) # Последний снимок пятидневной давности совпадает, а более старый — нет: # сравнивать обязаны с ПОСЛЕДНИМ, а не с любым из прошлых. _snapshot(db, ids["same_gap"], days_ago=9, price=_PRICE - 1, payload=p, seen_days_ago=9) _snapshot(db, ids["same_gap"], days_ago=5, price=_PRICE, payload=p, seen_days_ago=1) _snapshot(db, ids["price"], days_ago=1, price=4_900_000, payload=p, seen_days_ago=1) _snapshot(db, ids["payload"], days_ago=1, price=_PRICE, payload='{"v": 0}', seen_days_ago=1) _snapshot(db, ids["seen"], days_ago=1, price=_PRICE, payload=p, seen_days_ago=3) # last_seen_at тот же, но 10 суток > окна свежести: вчера снимок был активным. _snapshot(db, ids["active"], days_ago=1, price=_PRICE, payload=p, seen_days_ago=10) return ids def test_only_changed_or_new_sources_get_a_row_today() -> None: db = _live_session() try: ids = _seed(db) _run_writer(db) today = _today(db, list(ids.values())) by_tag = {tag: today.get(lsid) for tag, lsid in ids.items()} assert by_tag == { "same": None, "same_gap": None, "price": (_PRICE, True), "payload": (_PRICE, True), "seen": (_PRICE, True), "active": (_PRICE, False), "new": (_PRICE, True), } finally: db.rollback() db.close() def test_events_survive_change_only_writes() -> None: """event-diff видит только записанные сегодня строки — и не теряет ни одного события.""" db = _live_session() try: ids = _seed(db) _run_writer(db) db.execute(snap_mod._EVENT_DIFF_SQL).fetchall() rows = db.execute( text( """ SELECT listing_source_id, event_type, diff_percent FROM listing_source_events WHERE listing_source_id = ANY(:ids) """ ), {"ids": list(ids.values())}, ).fetchall() tag_of = {lsid: tag for tag, lsid in ids.items()} events = {(tag_of[int(r[0])], r[1], None if r[2] is None else float(r[2])) for r in rows} assert events == { ("price", "price_change", 2.0408), # round((5.0-4.9)/4.9*100, 4) ("payload", "edited", None), ("new", "first_seen", None), } finally: db.rollback() db.close() def test_rerun_same_day_compares_with_todays_row() -> None: """Второй прогон в те же сутки сравнивает с уже записанной сегодня строкой. price вернулся к вчерашней цене — сегодняшняя строка обязана перезаписаться (иначе на сегодня осталась бы цена, которой у источника больше нет); same изменился после первого прогона — строка появляется; seen не менялся после первого прогона — строка остаётся как была. """ db = _live_session() try: ids = _seed(db) _run_writer(db) db.execute( text("UPDATE listing_sources SET price_rub = 4900000 WHERE id = :id"), {"id": ids["price"]}, ) db.execute( text("UPDATE listing_sources SET price_rub = 5100000 WHERE id = :id"), {"id": ids["same"]}, ) _run_writer(db) today = _today(db, list(ids.values())) assert today.get(ids["price"]) == (4_900_000, True) assert today.get(ids["same"]) == (5_100_000, True) assert today.get(ids["seen"]) == (_PRICE, True) assert ids["same_gap"] not in today finally: db.rollback() db.close()