"""listings_snapshots writer — point-in-time observation per listing per day. Таблица listings_snapshots (016_listings_snapshots.sql): PRIMARY KEY (listing_id, snapshot_date) — максимум 1 snapshot в сутки. ON CONFLICT DO UPDATE — берём последний за день (перезаписываем при повторном run'е). Используется двумя путями: 1. SERP scrape (save_listings в base.py) — записывает цену + позицию в выдаче каждый раз когда listing появляется в поиске. 2. Detail backfill (save_detail_enrichment в cian_detail.py) — записывает цену из detail-страницы для listings которые раньше не имели snapshot'а. """ from __future__ import annotations import logging from datetime import date from sqlalchemy import text from sqlalchemy.orm import Session logger = logging.getLogger(__name__) def upsert_listing_snapshot( db: Session, *, listing_id: int, price_rub: int, price_per_m2: int | None = None, run_id: int | None = None, snapshot_date: date | None = None, position_in_serp: int | None = None, status: str | None = "active", ) -> None: """Записать / обновить snapshot для listing за текущий (или указанный) день. Идемпотентно: повторный вызов с теми же listing_id + snapshot_date перезаписывает строку (берём последние данные за день — ON CONFLICT DO UPDATE). Args: db: открытая SQLAlchemy session. Commit делает CALLER. listing_id: PK из таблицы listings. price_rub: текущая цена в рублях. price_per_m2: цена за кв.м (вычисляется в caller'е, None если нет). run_id: FK scrape_runs.id — с каким run'ом связан snapshot (None если backfill запускается вне run'а). snapshot_date: дата наблюдения; если None — используется CURRENT_DATE (по БД). position_in_serp: позиция в SERP (1..60+), только при SERP-scrape. status: 'active' / 'closed' / None. Default 'active' при обычном scrape. """ db.execute( text(""" INSERT INTO listings_snapshots (listing_id, snapshot_date, run_id, price_rub, price_per_m2, position_in_serp, status, observed_at) VALUES ( CAST(:lid AS bigint), COALESCE(CAST(:snap_date AS date), CURRENT_DATE), CAST(:run_id AS bigint), CAST(:price AS bigint), CAST(:ppm2 AS int), CAST(:pos AS int), :status, NOW() ) ON CONFLICT (listing_id, snapshot_date) DO UPDATE SET run_id = COALESCE(EXCLUDED.run_id, listings_snapshots.run_id), price_rub = EXCLUDED.price_rub, price_per_m2 = COALESCE(EXCLUDED.price_per_m2, listings_snapshots.price_per_m2), position_in_serp = COALESCE( EXCLUDED.position_in_serp, listings_snapshots.position_in_serp ), status = COALESCE(EXCLUDED.status, listings_snapshots.status), observed_at = EXCLUDED.observed_at """), { "lid": listing_id, "snap_date": snapshot_date, "run_id": run_id, "price": price_rub, "ppm2": price_per_m2, "pos": position_in_serp, "status": status, }, ) logger.debug( "upsert_listing_snapshot: listing_id=%s date=%s price=%s run_id=%s", listing_id, snapshot_date or "CURRENT_DATE", price_rub, run_id, )