Compare commits
No commits in common. "38e38e662cc385f302c66e3dc1843a012d8769f6" and "f5e2f5f5f3253da643d484297778e54d29a6d8a1" have entirely different histories.
38e38e662c
...
f5e2f5f5f3
3 changed files with 17 additions and 548 deletions
|
|
@ -276,43 +276,6 @@ _REVISIT_FLOOR_NULL_SEGMENT_FILTER = "\n AND l.listing_segment IS NULL"
|
||||||
CAP_MULT = 2
|
CAP_MULT = 2
|
||||||
|
|
||||||
|
|
||||||
def _revisit_floor_from_where_sql(staleness_column: str, segment_filter: str) -> str:
|
|
||||||
"""FROM..WHERE, общий для _build_revisit_floor_sql и _build_revisit_floor_pairs_count_sql.
|
|
||||||
|
|
||||||
Вынесено в отдельную функцию, а не продублировано, ЦЕЛЕНАПРАВЛЕННО: пункт 3 PR-A
|
|
||||||
(#2659 продолжение) требует, чтобы n_pairs в counters считался «по тому же срезу,
|
|
||||||
что даёт выборку для percentile_disc» -- дословно, а не приблизительно. Общий текст
|
|
||||||
гарантирует это по построению; раздельные копии SELECT count(*) и SELECT
|
|
||||||
percentile_disc(...) рано или поздно разошлись бы независимой правкой одной из
|
|
||||||
двух.
|
|
||||||
|
|
||||||
JOIN LATERAL — предшественник ищется ПО СТРОКЕ (per listing_source_id), не по
|
|
||||||
глобальной дате: `ORDER BY s.snapshot_date DESC LIMIT 1` берёт последний снимок
|
|
||||||
ЭТОГО listing_source_id не позже якоря (CURRENT_DATE - health_window_days),
|
|
||||||
независимо от того, писалась ли строка именно на дату якоря. Подробно — почему
|
|
||||||
именно так и что было раньше — см. докстринг _build_revisit_floor_sql.
|
|
||||||
"""
|
|
||||||
return f"""
|
|
||||||
FROM listings l
|
|
||||||
JOIN listing_sources ls
|
|
||||||
ON ls.listing_id = l.id
|
|
||||||
AND ls.ext_source = l.source
|
|
||||||
JOIN LATERAL (
|
|
||||||
SELECT s.last_seen_at
|
|
||||||
FROM listing_source_snapshots s
|
|
||||||
WHERE s.listing_source_id = ls.id
|
|
||||||
AND s.snapshot_date
|
|
||||||
<= CURRENT_DATE - CAST(:health_window_days AS integer)
|
|
||||||
ORDER BY s.snapshot_date DESC
|
|
||||||
LIMIT 1
|
|
||||||
) prev ON true
|
|
||||||
WHERE l.source = :listing_source
|
|
||||||
AND l.{staleness_column}
|
|
||||||
> NOW() - CAST(:health_window_days || ' days' AS interval)
|
|
||||||
AND l.{staleness_column} > prev.last_seen_at{segment_filter}
|
|
||||||
"""
|
|
||||||
|
|
||||||
|
|
||||||
def _build_revisit_floor_sql(
|
def _build_revisit_floor_sql(
|
||||||
staleness_column: str, *, with_segments: bool, null_segment_only: bool = False
|
staleness_column: str, *, with_segments: bool, null_segment_only: bool = False
|
||||||
) -> Any:
|
) -> Any:
|
||||||
|
|
@ -323,32 +286,6 @@ def _build_revisit_floor_sql(
|
||||||
Только строки, у которых свежесть реально сдвинулась, — то есть выжившие,
|
Только строки, у которых свежесть реально сдвинулась, — то есть выжившие,
|
||||||
а не «мы к ним не приходили».
|
а не «мы к ним не приходили».
|
||||||
|
|
||||||
«ПРЕДЫДУЩЕЕ НАБЛЮДЕНИЕ» -- ЧТО ЭТО ТОЧНО (PR-A, #2659 продолжение). Это ПОСЛЕДНЯЯ
|
|
||||||
строка listing_source_snapshots для данного listing_source_id, чей snapshot_date
|
|
||||||
не позже якоря (CURRENT_DATE - health_window_days) -- `ORDER BY snapshot_date DESC
|
|
||||||
LIMIT 1` в LATERAL-подзапросе _revisit_floor_from_where_sql. РАВЕНСТВО ПО ДАТЕ
|
|
||||||
(`prev.snapshot_date = <якорь>`) ЗДЕСЬ ЗАПРЕЩЕНО, и вот почему:
|
|
||||||
|
|
||||||
1) listing_source_snapshots переходит на модель «строка на изменение» (пишется
|
|
||||||
только когда значение отличается от предыдущего снимка, не ежедневно). В этой
|
|
||||||
модели equality-join по дате теряет подавляющее большинство пар: на 2026-08-20
|
|
||||||
«изменениями» являются 4 458 строк из 101 795 -- 95.6% пар для percentile_disc
|
|
||||||
пропадают, выборка схлопывается до n=3-4, а percentile_disc(0.99) на такой
|
|
||||||
выборке вырождается в максимум из трёх чисел, а не в реальный квантиль хвоста.
|
|
||||||
2) Дыры в ЕЖЕДНЕВНОЙ истории уже ломали equality-join и ДО перехода на change-only
|
|
||||||
модель -- на проде отмечены разрывы 03-14.06 (12 суток подряд), 03-04.07,
|
|
||||||
12.07, 26.07, 30-31.07, 01.08. В эти дни equality-join давал n_pairs=0,
|
|
||||||
floor_days=NULL, и пол молча не считался -- TTL оставался как задан, хотя
|
|
||||||
история для «ближайшего более раннего» снимка (см. комментарий про потолок
|
|
||||||
выше в модульном докстринге) в базе была, просто не РОВНО на эту дату.
|
|
||||||
|
|
||||||
Семантика при этом не меняется: между двумя изменениями last_seen_at по
|
|
||||||
определению постоянен (иначе строка была бы новым изменением), поэтому
|
|
||||||
«последний снимок не позже якоря» и «снимок ровно на дату якоря» при СПЛОШНОЙ
|
|
||||||
ежедневной истории дают одно и то же число -- разница проявляется только там,
|
|
||||||
где equality-join был неверен и раньше (гэпы) либо станет неверен при переходе
|
|
||||||
на change-only (почти повсеместно).
|
|
||||||
|
|
||||||
null_segment_only=True переопределяет with_segments -- IS NULL вместо ANY(:segments)
|
null_segment_only=True переопределяет with_segments -- IS NULL вместо ANY(:segments)
|
||||||
(ANY никогда не матчит NULL). На практике для null_segment_only-джобы этот запрос
|
(ANY никогда не матчит NULL). На практике для null_segment_only-джобы этот запрос
|
||||||
не строится вовсе (revisit_floor_quantile=0 -- см. модульный докстринг), но вариант
|
не строится вовсе (revisit_floor_quantile=0 -- см. модульный докстринг), но вариант
|
||||||
|
|
@ -370,35 +307,22 @@ def _build_revisit_floor_sql(
|
||||||
ORDER BY EXTRACT(epoch FROM (l.{staleness_column} - prev.last_seen_at))
|
ORDER BY EXTRACT(epoch FROM (l.{staleness_column} - prev.last_seen_at))
|
||||||
/ 86400.0
|
/ 86400.0
|
||||||
)
|
)
|
||||||
{_revisit_floor_from_where_sql(staleness_column, segment_filter)}
|
FROM listings l
|
||||||
"""
|
JOIN listing_sources ls
|
||||||
)
|
ON ls.listing_id = l.id
|
||||||
|
AND ls.ext_source = l.source
|
||||||
|
JOIN listing_source_snapshots prev
|
||||||
def _build_revisit_floor_pairs_count_sql(
|
ON prev.listing_source_id = ls.id
|
||||||
staleness_column: str, *, with_segments: bool, null_segment_only: bool = False
|
AND prev.snapshot_date = (
|
||||||
) -> Any:
|
SELECT max(snapshot_date)
|
||||||
"""count(*) пар, из которых percentile_disc в _build_revisit_floor_sql берёт квантиль.
|
FROM listing_source_snapshots
|
||||||
|
WHERE snapshot_date
|
||||||
Наблюдательность (PR-A, #2659 продолжение) -- НИЧЕГО не блокирует сейчас (гейт по
|
<= CURRENT_DATE - CAST(:health_window_days AS integer)
|
||||||
деградации выборки, если он когда-нибудь понадобится, -- отдельная задача, PR-B).
|
)
|
||||||
Пишется в counters["floor_n_pairs"] ДО того, как схлопнувшаяся выборка станет
|
WHERE l.source = :listing_source
|
||||||
видна только по повторению прод-инцидента, ради которого весь пол заведён (см.
|
AND l.{staleness_column}
|
||||||
ЗАМЕР НА ПРОДЕ 2026-08-09 в модульном докстринге).
|
> NOW() - CAST(:health_window_days || ' days' AS interval)
|
||||||
|
AND l.{staleness_column} > prev.last_seen_at{segment_filter}
|
||||||
Тот же срез, что и percentile_disc -- ОБЩАЯ функция _revisit_floor_from_where_sql,
|
|
||||||
не копия WHERE, см. её докстринг про то, почему это важно.
|
|
||||||
"""
|
|
||||||
if null_segment_only:
|
|
||||||
segment_filter = _REVISIT_FLOOR_NULL_SEGMENT_FILTER
|
|
||||||
elif with_segments:
|
|
||||||
segment_filter = _REVISIT_FLOOR_SEGMENT_FILTER
|
|
||||||
else:
|
|
||||||
segment_filter = ""
|
|
||||||
return text(
|
|
||||||
f"""
|
|
||||||
SELECT count(*)
|
|
||||||
{_revisit_floor_from_where_sql(staleness_column, segment_filter)}
|
|
||||||
"""
|
"""
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -575,10 +499,7 @@ def deactivate_stale_listings(
|
||||||
|
|
||||||
Returns {"deactivated": N} -- количество обновлённых строк (1:1 со снимками).
|
Returns {"deactivated": N} -- количество обновлённых строк (1:1 со снимками).
|
||||||
Если гейт не пропустил прогон: {"deactivated": 0, "confirmations": N,
|
Если гейт не пропустил прогон: {"deactivated": 0, "confirmations": N,
|
||||||
"skipped_unhealthy": 1} и НИ ОДНА строка не тронута. Если пол переобхода включён
|
"skipped_unhealthy": 1} и НИ ОДНА строка не тронута. Если пол переобхода поднял
|
||||||
(revisit_floor_quantile > 0): дополнительно {"floor_n_pairs": N} -- размер
|
|
||||||
выборки, из которой percentile_disc посчитал квантиль (наблюдательность, PR-A
|
|
||||||
#2659 продолжение; сейчас ничего не гейтит). Если пол при этом реально поднял
|
|
||||||
TTL: дополнительно {"revisit_floor_days": N, "ttl_days_effective": N}. Если пол
|
TTL: дополнительно {"revisit_floor_days": N, "ttl_days_effective": N}. Если пол
|
||||||
упёрся в потолок cap_mult: дополнительно {"ttl_floor_capped": 1,
|
упёрся в потолок cap_mult: дополнительно {"ttl_floor_capped": 1,
|
||||||
"ttl_days_floor_raw": N} -- N это то, во что пол поднял бы TTL БЕЗ потолка.
|
"ttl_days_floor_raw": N} -- N это то, во что пол поднял бы TTL БЕЗ потолка.
|
||||||
|
|
@ -711,21 +632,6 @@ def deactivate_stale_listings(
|
||||||
),
|
),
|
||||||
floor_params,
|
floor_params,
|
||||||
).scalar()
|
).scalar()
|
||||||
# Наблюдательность (PR-A, #2659 продолжение) -- та же выборка, что и
|
|
||||||
# percentile_disc выше (общая _revisit_floor_from_where_sql). Пока ничего
|
|
||||||
# не гейтит (это PR-B), но n_pairs обязан попасть в counters ДО того, как
|
|
||||||
# схлопнувшуюся выборку станет видно только по повторению прод-инцидента.
|
|
||||||
# Считается ВСЕГДА при включённом полу, в том числе когда floor_days
|
|
||||||
# получится NULL -- тогда n_pairs=0 и объясняет, почему пол не посчитался.
|
|
||||||
floor_n_pairs = db.execute(
|
|
||||||
_build_revisit_floor_pairs_count_sql(
|
|
||||||
staleness_column,
|
|
||||||
with_segments=segments is not None,
|
|
||||||
null_segment_only=null_segment_only,
|
|
||||||
),
|
|
||||||
floor_params,
|
|
||||||
).scalar()
|
|
||||||
counters["floor_n_pairs"] = int(floor_n_pairs or 0)
|
|
||||||
# NULL = истории снимков за окно нет вовсе (свежая БД, дыра в снимках).
|
# NULL = истории снимков за окно нет вовсе (свежая БД, дыра в снимках).
|
||||||
# Тогда пола нет и TTL остаётся как задан: выдумывать пол не из чего.
|
# Тогда пола нет и TTL остаётся как задан: выдумывать пол не из чего.
|
||||||
if floor_days is not None:
|
if floor_days is not None:
|
||||||
|
|
|
||||||
|
|
@ -109,12 +109,3 @@ tests/test_2992_upsert_unchanged_gate.py::test_listing_sources_unchanged_rescrap
|
||||||
# настоящую БД (в CI она есть, #2745), локально без TEST_DATABASE_URL пропускается. Строки
|
# настоящую БД (в CI она есть, #2745), локально без TEST_DATABASE_URL пропускается. Строки
|
||||||
# t3036-* тест удаляет в finally.
|
# t3036-* тест удаляет в finally.
|
||||||
tests/test_3036_detail_house_params_to_houses.py::test_live_fill_only_then_keep_then_unlinked_untouched
|
tests/test_3036_detail_house_params_to_houses.py::test_live_fill_only_then_keep_then_unlinked_untouched
|
||||||
|
|
||||||
# PR-A LATERAL-поиск предшественника в поле переобхода (#2659 продолжение) — тот же
|
|
||||||
# `_live_session()`. Три live-теста вставляют синтетическую историю снимков (в т.ч.
|
|
||||||
# с дырами) в ОДНОЙ транзакции и проверяют результат JOIN LATERAL/equality-join
|
|
||||||
# прямым SELECT'ом в той же транзакции — cleanup через rollback (без commit), явный
|
|
||||||
# DELETE не нужен. Без БД — skip; в CI Trade-In идут на postgres-сервисе (#2745).
|
|
||||||
tests/test_revisit_floor_lateral_lookup.py::test_lateral_ignores_gaps_and_matches_dense_daily_history
|
|
||||||
tests/test_revisit_floor_lateral_lookup.py::test_equality_join_loses_gapped_pairs_and_gives_a_smaller_floor
|
|
||||||
tests/test_revisit_floor_lateral_lookup.py::test_missing_history_before_anchor_still_yields_no_floor_via_lateral
|
|
||||||
|
|
|
||||||
|
|
@ -1,428 +0,0 @@
|
||||||
"""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.<col> > prev.last_seen_at) -- НЕ часть этой правки,
|
|
||||||
LATERAL меняет только то, КАК ищется prev, а не что с ним сравнивается."""
|
|
||||||
sql = str(task_mod._build_revisit_floor_sql("last_seen_at", with_segments=False).text)
|
|
||||||
assert "l.last_seen_at > prev.last_seen_at" in sql
|
|
||||||
Loading…
Add table
Reference in a new issue