"""LATERAL-поиск предшественника в поле переобхода (#2659 продолжение, PR-A). `_build_revisit_floor_sql` джойнил историю переобхода ПО РАВЕНСТВУ ДАТЫ: `prev.snapshot_date = (SELECT max(snapshot_date) FROM listing_source_snapshots WHERE snapshot_date <= CURRENT_DATE - health_window_days)`. Это ломается ДВАЖДЫ: 1. listing_source_snapshots переходит на модель «строка на изменение» (пишется только когда значение отличается от предыдущего снимка) -- в ней у подавляющего большинства listing_source_id на конкретную календарную дату строки просто нет. На 2026-08-20 «изменениями» являются 4 458 строк из 101 795 (95.6% пар теряются). 2. Дыры в истории ломали equality-join и раньше, до перехода на change-only (на проде: 03-14.06 -- 12 суток подряд, 03-04.07, 12.07, 26.07, 30-31.07, 01.08). Правка меняет equality-join на `JOIN LATERAL (... ORDER BY snapshot_date DESC LIMIT 1) prev ON true` -- предшественник ищется ПО СТРОКЕ (per listing_source_id), не по единой глобальной дате, и гэпы в истории для него прозрачны. Семантика не меняется: между изменениями last_seen_at по определению постоянен, поэтому «последний снимок не позже якоря» и «снимок ровно на дату якоря» при СПЛОШНОЙ ежедневной истории дают одно и то же число -- разница проявляется только там, где equality-join терял пары. Тесты ниже: - test_lateral_ignores_gaps_and_matches_dense_daily_history -- LATERAL на дырявой (change-only) истории даёт ТОТ ЖЕ floor_days, что дала бы сплошная суточная история для того же слушателя. - test_equality_join_loses_gapped_pairs_and_gives_a_smaller_floor -- на ОДНИХ И ТЕХ ЖЕ данных equality-join (воспроизведён буквально -- см. _OLD_EQUALITY_JOIN_FLOOR_SQL ниже, это ровно то условие, что было в _build_revisit_floor_sql до этого PR) теряет строки без снимка ровно на дату якоря и даёт МЕНЬШИЙ пол. - Две pure-SQL проверки без БД (форма LATERAL-запроса, "count(*)" из того же среза). Живая БД: см. `_live_session()` -- self-skip без реального Postgres, как соседние live-тесты в этом каталоге (test_2992_upsert_unchanged_gate.py и т.д.). В CI Trade-In есть postgres-сервис (ci-tradein.yml), локально без БД эти проверки skip'аются, а чисто-SQL тесты внизу файла бегут всегда. """ from __future__ import annotations import os import re import uuid from datetime import UTC, date, datetime, timedelta 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 deactivate_stale_avito as task_mod _HEALTH_WINDOW_DAYS = 5 def _live_session() -> Any | None: """Тот же контракт, что у соседних live-тестов (test_2992_upsert_unchanged_gate.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 # ── Воспроизведение СТАРОГО equality-join (буквально то, что было в # _build_revisit_floor_sql до этой правки) -- код удалён из app/, поэтому для # сравнения "до/после" запрос встроен сюда как литерал. Единственное упрощение: # вместо подзапроса `(SELECT max(snapshot_date) FROM listing_source_snapshots # WHERE snapshot_date <= CURRENT_DATE - :N)` (глобальный по ВСЕЙ таблице, не # скоупленный per-listing_source_id -- и в этом половина исходного бага) якорная # дата передаётся явно предвычисленным `CURRENT_DATE - :health_window_days`. # Само по себе равенство по дате -- ровно та семантика, что менялась в этом PR; # зависимость от глобального max() по ВСЕЙ таблице сделала бы этот тест хрупким # к данным других тестов в той же живой БД, никак не проверяя суть правки. _OLD_EQUALITY_JOIN_FLOOR_SQL = text( """ SELECT percentile_disc(CAST(:revisit_quantile AS double precision)) WITHIN GROUP ( ORDER BY EXTRACT(epoch FROM (l.last_seen_at - prev.last_seen_at)) / 86400.0 ) FROM listings l JOIN listing_sources ls ON ls.listing_id = l.id AND ls.ext_source = l.source JOIN listing_source_snapshots prev ON prev.listing_source_id = ls.id AND prev.snapshot_date = CURRENT_DATE - CAST(:health_window_days AS integer) WHERE l.source = :listing_source AND l.last_seen_at > NOW() - CAST(:health_window_days || ' days' AS interval) AND l.last_seen_at > prev.last_seen_at """ ) _OLD_EQUALITY_JOIN_COUNT_SQL = text( """ SELECT count(*) FROM listings l JOIN listing_sources ls ON ls.listing_id = l.id AND ls.ext_source = l.source JOIN listing_source_snapshots prev ON prev.listing_source_id = ls.id AND prev.snapshot_date = CURRENT_DATE - CAST(:health_window_days AS integer) WHERE l.source = :listing_source AND l.last_seen_at > NOW() - CAST(:health_window_days || ' days' AS interval) AND l.last_seen_at > prev.last_seen_at """ ) def _insert_pair( db: Any, *, source: str, ext_id: str, listing_last_seen_at: datetime, snapshots: list[tuple[date, datetime]], ) -> int: """Вставляет listings + listing_sources + N снимков listing_source_snapshots. БЕЗ commit -- вся синтетика живёт в ОДНОЙ транзакции теста и видна собственным же SELECT'ам (same-session read-your-writes), очистка -- rollback в конце теста, отдельного DELETE не нужно. """ 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 (:source, :url, :ext_id, :dedup_hash, 5000000, true, NOW(), :last_seen_at) RETURNING id """ ), { "source": source, "url": f"https://example.test/pr2659-lateral/{ext_id}", "ext_id": ext_id, "dedup_hash": f"pr2659-lateral-{ext_id}", "last_seen_at": listing_last_seen_at, }, ).scalar_one() listing_source_id = db.execute( text( """ INSERT INTO listing_sources (listing_id, ext_source, ext_id, confidence, matched_method) VALUES (:listing_id, :source, :ext_id, 1.0, 'test') RETURNING id """ ), {"listing_id": listing_id, "source": source, "ext_id": ext_id}, ).scalar_one() for snapshot_date, snapshot_last_seen_at in snapshots: db.execute( text( """ INSERT INTO listing_source_snapshots (listing_source_id, snapshot_date, is_active, last_seen_at) VALUES (:lsid, :snapshot_date, true, :last_seen_at) """ ), { "lsid": listing_source_id, "snapshot_date": snapshot_date, "last_seen_at": snapshot_last_seen_at, }, ) return int(listing_source_id) def _floor(db: Any, *, source: str, quantile: float = 1.0) -> float | None: result = db.execute( task_mod._build_revisit_floor_sql("last_seen_at", with_segments=False), { "listing_source": source, "health_window_days": _HEALTH_WINDOW_DAYS, "revisit_quantile": quantile, }, ).scalar() return float(result) if result is not None else None def _n_pairs(db: Any, *, source: str) -> int: result = db.execute( task_mod._build_revisit_floor_pairs_count_sql("last_seen_at", with_segments=False), {"listing_source": source, "health_window_days": _HEALTH_WINDOW_DAYS}, ).scalar() return int(result or 0) def _old_floor(db: Any, *, source: str, quantile: float = 1.0) -> float | None: result = db.execute( _OLD_EQUALITY_JOIN_FLOOR_SQL, { "listing_source": source, "health_window_days": _HEALTH_WINDOW_DAYS, "revisit_quantile": quantile, }, ).scalar() return float(result) if result is not None else None def _old_n_pairs(db: Any, *, source: str) -> int: result = db.execute( _OLD_EQUALITY_JOIN_COUNT_SQL, {"listing_source": source, "health_window_days": _HEALTH_WINDOW_DAYS}, ).scalar() return int(result or 0) # ── Живые тесты (реальный Postgres, транзакция + rollback) ──────────────────── @pytest.mark.skipif(_live_session() is None, reason="no reachable Postgres test DB") def test_lateral_ignores_gaps_and_matches_dense_daily_history() -> None: """Дырявая (change-only) история даёт ТОТ ЖЕ floor, что дала бы сплошная суточная. Между изменениями last_seen_at по определению постоянен -- поэтому у одного и того же listing_source_id снимок «ровно на дату якоря» и «последний снимок не позже якоря» несут ОДНО И ТО ЖЕ last_seen_at, если между ними ничего не менялось. Сначала считаем пол на СПЛОШНОЙ суточной истории (якорная дата присутствует), затем удаляем все снимки, кроме самого раннего (change-only модель — сутки без изменений не пишутся), и проверяем, что LATERAL даёт то же число. """ db = _live_session() source = f"zzzlat_dense_{uuid.uuid4().hex[:8]}" now = datetime.now(UTC) anchor = date.today() - timedelta(days=_HEALTH_WINDOW_DAYS) # last_seen_at не меняется все эти сутки -- то самое "между изменениями постоянен". unchanged_last_seen_at = now - timedelta(days=_HEALTH_WINDOW_DAYS + 3) try: db.execute(text("BEGIN")) lsid = _insert_pair( db, source=source, ext_id="dense", listing_last_seen_at=now, snapshots=[ (anchor - timedelta(days=3), unchanged_last_seen_at), (anchor - timedelta(days=2), unchanged_last_seen_at), (anchor - timedelta(days=1), unchanged_last_seen_at), (anchor, unchanged_last_seen_at), # якорная дата ПРИСУТСТВУЕТ ], ) floor_dense = _floor(db, source=source) assert floor_dense is not None assert _n_pairs(db, source=source) == 1 # Change-only модель: сутки без изменений не пишутся -- удаляем всё, кроме # самого раннего снимка (якорная дата больше НЕ присутствует ровно). db.execute( text( "DELETE FROM listing_source_snapshots " "WHERE listing_source_id = :lsid AND snapshot_date > :keep_date" ), {"lsid": lsid, "keep_date": anchor - timedelta(days=3)}, ) floor_sparse = _floor(db, source=source) assert floor_sparse is not None assert _n_pairs(db, source=source) == 1 assert floor_sparse == pytest.approx(floor_dense, abs=0.01), ( f"LATERAL обязан игнорировать гэп: сплошная история дала {floor_dense}, " f"дырявая -- {floor_sparse}, а между изменениями last_seen_at постоянен" ) # Возраст известен точно (постоянный last_seen_at, N+3 суток разрыва) -- # пиним не только "совпадают", но и КАКОЕ именно число. assert floor_dense == pytest.approx(_HEALTH_WINDOW_DAYS + 3, abs=0.01) finally: db.rollback() db.close() @pytest.mark.skipif(_live_session() is None, reason="no reachable Postgres test DB") def test_equality_join_loses_gapped_pairs_and_gives_a_smaller_floor() -> None: """На ОДНИХ И ТЕХ ЖЕ данных equality-join теряет пары и даёт МЕНЬШИЙ пол -- это и есть суть правки PR-A: 3 строки, только у одной снимок ровно на дату якоря. Group A -- снимок РОВНО на дату якоря, возраст ~5.5 сут (выживает в ОБЕИХ моделях). Group B, C -- снимок ТОЛЬКО раньше якоря (change-only гэп ровно на дату якоря), возраст 65 и 30 сут -- выживают ТОЛЬКО в LATERAL. revisit_quantile=1.0 (максимум) -- детерминированно и без интерполяции percentile_disc на маленькой выборке: LATERAL обязан дать max(5.5, 65, 30) = 65, equality-join -- только 5.5 (единственная сохранившаяся пара). """ db = _live_session() source = f"zzzlat_gap_{uuid.uuid4().hex[:8]}" now = datetime.now(UTC) anchor = date.today() - timedelta(days=_HEALTH_WINDOW_DAYS) try: db.execute(text("BEGIN")) # Group A: снимок ровно на дату якоря -- переживает equality-join. _insert_pair( db, source=source, ext_id="group-a-on-anchor", listing_last_seen_at=now, snapshots=[(anchor, now - timedelta(days=5, hours=12))], # возраст 5.5 сут ) # Group B: снимок только на anchor-60 -- дыра ровно на дату якоря # (типичный change-only разрыв: ничего не менялось 55+ суток подряд). _insert_pair( db, source=source, ext_id="group-b-gap-60d", listing_last_seen_at=now, snapshots=[(anchor - timedelta(days=60), now - timedelta(days=65))], ) # Group C: тот же класс гэпа, гэп короче (30 сут) -- контроль, что LATERAL # берёт максимум по ВСЕЙ выборке, а не просто "последнюю вставленную пару". _insert_pair( db, source=source, ext_id="group-c-gap-25d", listing_last_seen_at=now, snapshots=[(anchor - timedelta(days=25), now - timedelta(days=30))], ) lateral_floor = _floor(db, source=source) lateral_n_pairs = _n_pairs(db, source=source) old_floor = _old_floor(db, source=source) old_n_pairs = _old_n_pairs(db, source=source) assert lateral_n_pairs == 3, "LATERAL обязан видеть все 3 пары, гэпы не теряют строк" assert old_n_pairs == 1, "equality-join теряет Group B и C -- у них нет снимка на якоре" assert lateral_floor == pytest.approx(65.0, abs=0.01), ( f"LATERAL max должен взять Group B (65 сут), получили {lateral_floor}" ) assert old_floor == pytest.approx(5.5, abs=0.01), ( f"equality-join должен остаться только с Group A (5.5 сут), получили {old_floor}" ) assert lateral_floor > old_floor, ( "равенство по дате даёт МЕНЬШИЙ пол на тех же данных -- это и есть регрессия, " "которую эта правка чинит" ) finally: db.rollback() db.close() @pytest.mark.skipif(_live_session() is None, reason="no reachable Postgres test DB") def test_missing_history_before_anchor_still_yields_no_floor_via_lateral() -> None: """Контроль: если у listing_source_id вообще нет снимка не позже якоря (совсем свежая строка), LATERAL закономерно не находит пару -- НЕ ломается на NULL/пусто, а просто исключает строку (как и раньше исключал equality-join, если совсем нет истории). Это НЕ регрессия, а ожидаемое поведение при отсутствии данных.""" db = _live_session() source = f"zzzlat_nohist_{uuid.uuid4().hex[:8]}" now = datetime.now(UTC) anchor = date.today() - timedelta(days=_HEALTH_WINDOW_DAYS) try: db.execute(text("BEGIN")) _insert_pair( db, source=source, ext_id="only-future-snapshot", listing_last_seen_at=now, # Снимок ЕСТЬ, но он ПОЗЖЕ якоря -- LATERAL (snapshot_date <= якорь) # его не видит, что и требуется: свежая строка без прошлого не должна # выдумывать пол из снимка, который сам моложе окна здоровья. snapshots=[(anchor + timedelta(days=1), now - timedelta(days=1))], ) assert _floor(db, source=source) is None assert _n_pairs(db, source=source) == 0 finally: db.rollback() db.close() # ── Pure-SQL проверки (без БД) -- форма LATERAL-запроса и общего среза ──────── def test_lateral_sql_has_no_equality_join_on_snapshot_date() -> None: """Ключевой негативный инвариант правки: РАВЕНСТВО по дате запрещено.""" sql = str(task_mod._build_revisit_floor_sql("last_seen_at", with_segments=False).text) assert "JOIN LATERAL" in sql assert "ORDER BY s.snapshot_date DESC" in sql assert "LIMIT 1" in sql assert not re.search(r"prev\.snapshot_date\s*=", sql), ( "equality-join по snapshot_date должен быть полностью удалён из пола переобхода" ) def test_lateral_sql_still_psycopg_v3_safe() -> None: sql = str(task_mod._build_revisit_floor_sql("last_seen_at", with_segments=True).text) assert "CAST(:revisit_quantile AS double precision)" in sql assert "CAST(:health_window_days AS integer)" in sql assert not re.search(r":\w+::", sql) def test_pairs_count_sql_shares_the_exact_same_slice_as_the_floor_sql() -> None: """count(*) и percentile_disc обязаны идти по ОДНОМУ И ТОМУ ЖЕ FROM..WHERE -- иначе floor_n_pairs в counters не описывает реальную выборку percentile_disc.""" floor_sql = str(task_mod._build_revisit_floor_sql("last_seen_at", with_segments=True).text) count_sql = str( task_mod._build_revisit_floor_pairs_count_sql("last_seen_at", with_segments=True).text ) floor_from_where = floor_sql.split("FROM listings l", 1)[1] count_from_where = count_sql.split("FROM listings l", 1)[1] assert floor_from_where == count_from_where, ( "FROM..WHERE percentile_disc-запроса и count(*)-запроса разошлись -- " "floor_n_pairs больше не описывает реальную выборку пола" ) def test_comparison_predicate_against_prev_last_seen_at_is_untouched() -> None: """Предикат сравнения (l.