Some checks failed
CI Trade-In / changes (pull_request) Successful in 26s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 31s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Failing after 7m37s
listing_source_snapshot каждую ночь копировал listing_sources целиком: прод 14-17.09 — 290-293 тыс. строк в сутки, из них отличались от предыдущего снимка 7.7-11.4 тыс. (3-4 %). Решение владельца 2026-08-23 — строка на изменение; читатель (пол переобхода, PR-A #3056) и рельсы объёма (PR-B #3066) уже смержены. Строка пишется, если снимка ещё нет или (price_rub, is_active, last_seen_at, payload_hash) отличается от последнего снимка источника, включая сегодняшний (повторный прогон в те же сутки перезаписывает строку, только если значение снова сдвинулось). Состояние на дату D = последняя строка с snapshot_date <= D — для всех четырёх колонок то же, что дала бы суточная копия, поэтому пол переобхода и события не меняются. Старые суточные снимки не трогаются. Замер чтением на проде 17.09: отбор с условием — 1.0 с на 293 602 источника (бюджет 900 с). Миграция 325 — только COMMENT на таблицу и view: прежние «ежедневный снимок»/«на дату снимка» стали бы неправдой. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
246 lines
9.3 KiB
Python
246 lines
9.3 KiB
Python
"""#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()
|