fix(tradein): писатели наконец пишут то, что обещает схема — фото подсказок, статус «снято», события объявлений (#2674) #2682

Merged
bot-backend merged 3 commits from fix/2674-writers-honor-schema into main 2026-08-05 22:12:30 +00:00
9 changed files with 766 additions and 80 deletions

View file

@ -345,26 +345,37 @@ def save_imv_result(db: Session, house_id: int, params: dict, result: IMVEvaluat
)
# 3. Suggestions
# #2674: до этого фикса в INSERT не входили image_link + area_m2/rooms/floor/
# total_floors — колонки есть с миграции 064, но писатель их не заполнял
# (25 055 строк на проде с NULL во всех пяти). Ссылка на фото приходит в
# suggestions.items[].imageLink, метрики квартиры парсятся из title.
for sug in result.suggestions:
db.execute(
text("""
INSERT INTO house_suggestions (
house_id, ext_item_id, title, address, price_rub,
area_m2, rooms, floor, total_floors,
exposure_days, publish_date,
item_link, metro_name, metro_distance,
item_link, image_link, metro_name, metro_distance,
has_good_price_badge, raw_payload, fetched_at
) VALUES (
:hid, :ext, :title, :addr, :price,
CAST(:area AS numeric), :rooms, :floor, :total_floors,
:exp, :pdate,
:link, :mname, :mdist,
:link, :img, :mname, :mdist,
:gpb, CAST(:raw AS jsonb), NOW()
)
ON CONFLICT (house_id, ext_item_id) DO UPDATE SET
title = EXCLUDED.title,
price_rub = EXCLUDED.price_rub,
area_m2 = EXCLUDED.area_m2,
rooms = EXCLUDED.rooms,
floor = EXCLUDED.floor,
total_floors = EXCLUDED.total_floors,
exposure_days = EXCLUDED.exposure_days,
publish_date = EXCLUDED.publish_date,
item_link = EXCLUDED.item_link,
image_link = EXCLUDED.image_link,
metro_name = EXCLUDED.metro_name,
metro_distance = EXCLUDED.metro_distance,
has_good_price_badge = EXCLUDED.has_good_price_badge,
@ -377,9 +388,14 @@ def save_imv_result(db: Session, house_id: int, params: dict, result: IMVEvaluat
"title": sug.title,
"addr": sug.address,
"price": sug.price_rub,
"area": sug.area_m2,
"rooms": sug.rooms,
"floor": sug.floor,
"total_floors": sug.total_floors,
"exp": sug.exposure_days,
"pdate": sug.publish_date,
"link": sug.item_url,
"img": sug.image_link,
"mname": sug.metro_name,
"mdist": sug.metro_distance,
"gpb": sug.has_good_price_badge,

View file

@ -47,6 +47,7 @@ from scraper_kit.providers.avito.detail import (
save_detail_enrichment,
)
from scraper_kit.providers.avito.serp import AvitoScraper
from scraper_kit.snapshot_writer import upsert_listing_snapshot
from sqlalchemy import text
from sqlalchemy.orm import Session
@ -241,7 +242,7 @@ async def run_avito_detail_backfill(
text(
"""
WITH ekb AS (
SELECT id, source_url, 'ekb' AS city_scope
SELECT id, source_url, price_rub, 'ekb' AS city_scope
FROM listings
WHERE source = 'avito'
AND detail_enriched_at IS NULL
@ -254,7 +255,7 @@ async def run_avito_detail_backfill(
LIMIT CAST(:batch_size AS int)
),
oblast AS (
SELECT id, source_url, 'oblast' AS city_scope
SELECT id, source_url, price_rub, 'oblast' AS city_scope
FROM listings
WHERE source = 'avito'
AND detail_enriched_at IS NULL
@ -264,9 +265,9 @@ async def run_avito_detail_backfill(
ORDER BY (lat IS NULL) DESC, scraped_at DESC NULLS LAST
LIMIT CAST(:oblast_batch_size AS int)
)
SELECT id, source_url, city_scope FROM ekb
SELECT id, source_url, price_rub, city_scope FROM ekb
UNION ALL
SELECT id, source_url, city_scope FROM oblast
SELECT id, source_url, price_rub, city_scope FROM oblast
"""
),
{
@ -442,6 +443,21 @@ async def run_avito_detail_backfill(
text("UPDATE listings SET is_active = FALSE WHERE id = :id"),
{"id": row["id"]},
)
# #2674: 404 с площадки — самый достоверный сигнал снятия,
# фиксируем его в дневной истории (listings_snapshots.status
# был константой 'active' у всех строк, 394 299). Тот же
# SAVEPOINT, что и UPDATE флага: снимок без флага (или
# наоборот) невозможен. price_rub из snapshot-SELECT —
# .get() консервативен ради mock-снапшотов старых тестов.
gone_price = row.get("price_rub")
if gone_price is not None:
upsert_listing_snapshot(
db,
listing_id=row["id"],
price_rub=gone_price,
run_id=run_id,
status="closed",
)
except Exception:
logger.warning(
"avito_detail_backfill: run_id=%d failed to mark listing %s "

View file

@ -9,6 +9,9 @@
TTL=30. novostroyki (9659 активных первичных строк) и NULL-сегмент не трогаем.
- avito: все сегменты (segments=None), TTL=10 дней -- поведение без изменений.
- Строки НЕ удаляются -- история нужна для бэктеста (#667).
- #2674: деактивация в той же транзакции пишет снимок listings_snapshots со статусом
'stale' за текущую дату -- «мы N суток не видели». Жёсткое 'closed' (площадка
ответила 404) пишет только avito_detail_backfill: смешивать факт с догадкой дорого.
Задача синхронная (DB-only, никаких внешних HTTP-вызовов) -- запускается kit-scheduler'ом
через product_handlers._job_deactivate_stale (wildcard-handler deactivate_stale_*),
@ -41,21 +44,62 @@ logger = logging.getLogger(__name__)
# только реальный скрейп).
_ALLOWED_STALENESS_COLUMNS = frozenset({"last_seen_at", "scraped_at"})
# ── Снимок «протухло» в дневной истории (#2674) ───────────────────────────────
# listings_snapshots.status до этого фикса был константой 'active' у всех строк
# (394 299 на момент находки) — оба места вызова upsert_listing_snapshot передавали
# литерал 'active', и это честно: там объявление ДЕЙСТВИТЕЛЬНО видели. А деактивация
# TTL-задачей не оставляла в истории вообще никакого следа. Из-за этого дата снятия
# объявления (лучший доступный сигнал «скорее всего продано») не запрашивалась из
# истории, а восстанавливалась на глаз: последний показ + предполагаемый срок жизни.
#
# ПОЧЕМУ 'stale', А НЕ 'closed'. Эта задача НЕ знает, что объявление снято, — она
# знает только, что МЫ его N суток не видели, а это разные факты, когда TTL короче
# простоя обхода. Замер: прогон по домклику 02.08 снял 6131 объявление за раз (TTL
# 14 суток против 12 суток простоя обхода) — с общим статусом это были бы 6131
# фальшивая «дата продажи» одной датой. Продукт про цены, смешивать факт с догадкой
# дорого. Поэтому:
# 'closed' — только путь 404: площадка ответила «нет» (avito_detail_backfill);
# 'stale' — этот путь: «мы N суток не смотрели».
# Дата всё равно фиксируется, но читатель отличает одно от другого. Ограничения
# CHECK на колонке нет (проверено на проде), миграция не нужна — только COMMENT.
#
# Пишем снимок в ТОЙ ЖЕ транзакции, что и UPDATE флага: деактивация без снимка (или
# наоборот) невозможна по построению — один statement, data-modifying CTE.
# 1:1 по строкам: `stale` возвращает уникальные listings.id (PK), каждая даёт ровно
# одну затронутую строку listings_snapshots (INSERT либо DO UPDATE — оба считаются
# в rowcount), поэтому rowcount statement'а по-прежнему равен числу деактивированных.
# price_rub берём из listings (NOT NULL в схеме) — это последняя известная цена.
# ON CONFLICT: если снимок за сегодня уже есть (объявление видели активным утром,
# а вечером сработал TTL) — только переводим статус в 'stale', цену не переписываем.
_STALE_SNAPSHOT_TAIL = """
INSERT INTO listings_snapshots
(listing_id, snapshot_date, run_id, price_rub, status, observed_at)
SELECT id, CURRENT_DATE, CAST(:run_id AS bigint), price_rub, 'stale', NOW()
FROM stale
ON CONFLICT (listing_id, snapshot_date) DO UPDATE SET
status = 'stale',
observed_at = EXCLUDED.observed_at
"""
def _build_all_segments_sql(staleness_column: str) -> Any:
"""UPDATE без фильтра по сегменту: все сегменты для данного source.
staleness_column уже прошёл whitelist-проверку в deactivate_stale_listings,
поэтому f-string-подстановка имени колонки безопасна. Значения (:listing_source,
:ttl_days) остаются param-binding psycopg v3 safe (никаких :param::type).
:ttl_days, :run_id) остаются param-binding psycopg v3 safe (никаких :param::type).
"""
return text(
f"""
UPDATE listings
SET is_active = false
WHERE source = :listing_source
AND is_active = true
AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval)
WITH stale AS (
UPDATE listings
SET is_active = false
WHERE source = :listing_source
AND is_active = true
AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval)
RETURNING id, price_rub
)
{_STALE_SNAPSHOT_TAIL}
"""
)
@ -68,12 +112,16 @@ def _build_segments_sql(staleness_column: str) -> Any:
"""
return text(
f"""
UPDATE listings
SET is_active = false
WHERE source = :listing_source
AND is_active = true
AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval)
AND listing_segment = ANY(CAST(:segments AS text[]))
WITH stale AS (
UPDATE listings
SET is_active = false
WHERE source = :listing_source
AND is_active = true
AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval)
AND listing_segment = ANY(CAST(:segments AS text[]))
RETURNING id, price_rub
)
{_STALE_SNAPSHOT_TAIL}
"""
)
@ -111,9 +159,10 @@ def deactivate_stale_listings(
свежесть = scraped_at (двигается только реальным скрейпом).
Sync (вызывается scheduler-триггером в executor, как snapshot_listing_sources).
Один UPDATE в транзакции. Финализирует scrape_runs (mark_done / mark_failed).
Один statement в транзакции: UPDATE флага + снимок 'stale' в listings_snapshots
(data-modifying CTE, #2674). Финализирует scrape_runs (mark_done / mark_failed).
Returns {"deactivated": N} -- количество обновлённых строк.
Returns {"deactivated": N} -- количество обновлённых строк (1:1 со снимками).
Raises:
ValueError: если staleness_column не входит в whitelist (проверка ДО SQL,
@ -140,10 +189,15 @@ def deactivate_stale_listings(
"listing_source": listing_source,
"ttl_days": ttl_days,
"segments": segments,
"run_id": run_id,
}
result = db.execute(_build_segments_sql(staleness_column), params)
else:
params = {"listing_source": listing_source, "ttl_days": ttl_days}
params = {
"listing_source": listing_source,
"ttl_days": ttl_days,
"run_id": run_id,
}
result = db.execute(_build_all_segments_sql(staleness_column), params)
counters["deactivated"] = result.rowcount or 0

View file

@ -104,13 +104,51 @@ _SNAPSHOT_SQL = text(
"""
)
# ── Event diff: price_change ──────────────────────────────────────────────────
# Для каждого источника сравниваем сегодняшнюю цену (snapshot_date = CURRENT_DATE) с
# самым свежим ПРЕДЫДУЩИМ снимком (snapshot_date < CURRENT_DATE). Если цена изменилась
# (обе NOT NULL, old <> 0) — пишем price_change.
# ── Event diff: три выводимых типа событий из пяти в схеме ────────────────────
# Для каждого источника сравниваем сегодняшний снимок (snapshot_date = CURRENT_DATE) с
# самым свежим ПРЕДЫДУЩИМ (snapshot_date < CURRENT_DATE).
# today — снимок за сегодня (только что записан _SNAPSHOT_SQL, в той же транзакции).
# p — последний снимок строго ДО сегодня, per-row LATERAL point-lookup (#2607).
#
# #2674: схема (079) знает пять типов событий, писатель умел один — price_change,
# 8288 строк. Дописаны два:
# edited — payload_hash изменился, а цена нет (изменение цены уже описано
# отдельным событием price_change — дублировать его как «редактирование»
# значило бы считать одно изменение дважды). Прошлый хеш обязан быть
# непустым: md5(NULL) = NULL, и «payload появился впервые» — это не
# правка, а первое наблюдение;
# first_seen — предыдущего снимка нет вовсе (LEFT JOIN LATERAL даёт p.* = NULL).
#
# delisted и relisted НЕ ПИШУТСЯ НАМЕРЕННО — они НЕ ВЫВОДИМЫ из наших данных.
# is_active в снимке — derived-признак «last_seen_at свежее FRESHNESS_WINDOW_DAYS»,
# то есть «мы видели», а не «объявление есть на площадке». При покрытии обхода 10-35%
# такой переход рождается тем, что скрейпер СНОВА ДОШЁЛ до источника, а не тем, что
# объявление вернулось/ушло. Контрольная группа в наших же данных (14-18.07):
# domklik, покрытие 99.9-100%: снятий 1/2/0/2/4 в сутки, возвратов — РОВНО 0 все дни;
# yandex, покрытие 34-43%: снятий 343-433 в сутки, возвратов до 155.
# Тот же обход, тот же день — разница только в покрытии. Отсюда же всплески:
# avito 13.07 (день остановки обхода) — 3023 «снятия» за сутки против контрольной
# ставки 1-4, точность события ≈4%; 4705 «возвратов» из 5493 за 12 дней (86%) — это
# два дня после возобновления обхода 2-3.08.
# Сузить окно свежести НЕ поможет — станет хуже (больше флапаний); окно шире
# максимального интервала повторного визита обессмысливает само событие.
# Честный ответ схеме — не писать эти два типа, а не наполнять журнал догадками.
# Единственный жёсткий сигнал снятия — 404 при поштучном обходе, он пишется в
# listings_snapshots.status='closed' (avito_detail_backfill).
#
# Оставшиеся три события утверждают факты о НАШИХ СОБСТВЕННЫХ строках («появился новый
# источник», «хеш изменился при той же цене», «цена другая»), а не о поведении площадки.
#
# JOIN → LEFT JOIN LATERAL: без LEFT источники без предыдущего снимка отбрасывались
# join'ом, поэтому first_seen был недостижим по построению. План #2607 не меняется —
# LEFT JOIN LATERAL так же форсирует per-row индексный point-lookup по
# idx_lss_source_date, просто не отбрасывает строку при отсутствии предыдущей.
#
# Ветки разворачиваются CROSS JOIN LATERAL (VALUES ...) — одна строка сравнения даёт
# до трёх строк-кандидатов, из которых WHERE e.fires оставляет сработавшие. Это
# по-прежнему ОДИН set-based statement (никакого Python-цикла), просто три предиката
# вместо одного.
#
# #2607: раньше `p` был отдельным CTE `DISTINCT ON (listing_source_id) ... FROM
# listing_source_snapshots WHERE snapshot_date < CURRENT_DATE` и джойнился обычным JOIN.
# Планировщик оценивает `today` в 1 строку (свежевставленные в этой же транзакции строки
@ -127,38 +165,79 @@ _SNAPSHOT_SQL = text(
#
# Полностью set-based: один INSERT … SELECT по всем источникам, без Python-цикла (LATERAL
# — это внутренний план Postgres, не Python-итерация).
# change_time = now() детерминирует UNIQUE(listing_source_id, change_time, event_type)
# в пределах прогона → ON CONFLICT DO NOTHING делает писатель идемпотентным.
#
# change_time = date_trunc('day', now()), а НЕ now() (#2674): с now() уникальность
# UNIQUE(listing_source_id, change_time, event_type) работала только ВНУТРИ прогона —
# второй прогон в те же сутки перезаписывал сегодняшний снимок, предикаты срабатывали
# заново с другим временем и давали дубли (2 августа таких прогонов было два).
# Суточная гранулярность честнее для суточного же сравнения и включает заявленную
# идемпотентность: ON CONFLICT DO NOTHING теперь действительно гасит повтор за день.
#
# NULLIF(p.price_rub, 0) в diff_percent обязателен: выражения VALUES вычисляются ДО
# фильтра `WHERE e.fires`, поэтому предикат "p.price_rub <> 0" от деления на ноль уже
# не спасает — без NULLIF первый же источник с нулевой прошлой ценой уронил бы весь
# прогон. Результат при этом тот же: строка с NULL-диффом не проходит e.fires.
#
# Внешний SELECT над data-modifying CTE считает вставленное ПО ТИПАМ (RETURNING отдаёт
# только реально вставленные строки, не съеденные ON CONFLICT), сразу в виде ключей
# счётчиков `<event_type>_events` — писатель получает готовый dict без Python-агрегации.
# Ровно этот счётчик и показал бы четыре нуля из пяти, если бы существовал раньше.
_EVENT_DIFF_SQL = text(
"""
WITH today AS (
SELECT listing_source_id, price_rub
SELECT listing_source_id, price_rub, payload_hash
FROM listing_source_snapshots
WHERE snapshot_date = CURRENT_DATE
),
inserted AS (
INSERT INTO listing_source_events (
listing_source_id, change_time, event_type, price_rub, diff_percent
)
SELECT
t.listing_source_id,
date_trunc('day', now()),
e.event_type,
t.price_rub,
e.diff_percent
FROM today t
LEFT JOIN LATERAL (
SELECT s.snapshot_date, s.price_rub, s.payload_hash
FROM listing_source_snapshots s
WHERE s.listing_source_id = t.listing_source_id
AND s.snapshot_date < CURRENT_DATE
ORDER BY s.snapshot_date DESC
LIMIT 1
) p ON true
CROSS JOIN LATERAL (VALUES
(
'first_seen',
NULL::numeric,
p.snapshot_date IS NULL
),
(
'price_change',
round((t.price_rub - p.price_rub)::numeric
/ NULLIF(p.price_rub, 0) * 100, 4),
t.price_rub IS NOT NULL
AND p.price_rub IS NOT NULL
AND p.price_rub <> 0
AND t.price_rub <> p.price_rub
),
(
'edited',
NULL::numeric,
p.payload_hash IS NOT NULL
AND t.payload_hash IS DISTINCT FROM p.payload_hash
AND t.price_rub IS NOT DISTINCT FROM p.price_rub
)
) AS e(event_type, diff_percent, fires)
WHERE e.fires
ON CONFLICT (listing_source_id, change_time, event_type) DO NOTHING
RETURNING event_type
)
INSERT INTO listing_source_events (
listing_source_id, change_time, event_type, price_rub, diff_percent
)
SELECT
t.listing_source_id,
now(),
'price_change',
t.price_rub,
round((t.price_rub - p.price_rub)::numeric / p.price_rub * 100, 4)
FROM today t
JOIN LATERAL (
SELECT s.price_rub
FROM listing_source_snapshots s
WHERE s.listing_source_id = t.listing_source_id
AND s.snapshot_date < CURRENT_DATE
ORDER BY s.snapshot_date DESC
LIMIT 1
) p ON true
WHERE t.price_rub IS NOT NULL
AND p.price_rub IS NOT NULL
AND p.price_rub <> 0
AND t.price_rub <> p.price_rub
ON CONFLICT (listing_source_id, change_time, event_type) DO NOTHING
SELECT event_type || '_events' AS counter_key, count(*) AS n
FROM inserted
GROUP BY 1
"""
)
@ -166,12 +245,14 @@ _EVENT_DIFF_SQL = text(
def snapshot_listing_sources(
db: Session, run_id: int, params: dict[str, Any] | None = None
) -> dict[str, int]:
"""Записать дневной снимок listing_sources + price_change-события.
"""Записать дневной снимок listing_sources + события изменений.
Sync (вызывается scheduler-триггером в executor, как import_rosreestr_dkp).
Два set-based statement'а в одной транзакции:
1. upsert снимка на (listing_source_id, CURRENT_DATE) last-write-wins.
2. diff сегодняшней цены против последнего предыдущего снимка price_change-события.
2. diff сегодняшнего снимка против последнего предыдущего три события,
выводимые из наших данных (#2674). delisted/relisted схема разрешает, но
они НЕ выводимы при покрытии обхода 10-35% см. _EVENT_DIFF_SQL.
Params (из default_params jsonb в scrape_schedules, #2607):
budget_sec: float SET LOCAL statement_timeout на транзакцию (default 900,
@ -183,11 +264,18 @@ def snapshot_listing_sources(
Финализирует scrape_runs (mark_done / mark_failed) и пишет counters.
Returns {"snapshotted": N, "price_change_events": M}.
Returns {"snapshotted": N, "<event_type>_events": M} по счётчику на каждый из
трёх пишущихся типов, всегда все три ключа (тип, который за прогон не сработал
ни разу, честно показывает 0, а не пропадает из counters).
"""
params = params or {}
budget_sec = _clamp_budget_sec(params.get("budget_sec", DEFAULT_BUDGET_SEC))
counters: dict[str, int] = {"snapshotted": 0, "price_change_events": 0}
counters: dict[str, int] = {
"snapshotted": 0,
"price_change_events": 0,
"edited_events": 0,
"first_seen_events": 0,
}
try:
# statement_timeout НЕ принимает bind-параметр ($1/:name) — синтаксис Postgres SET
# запрещает placeholder на этом месте (проверено вживую на проде: "syntax error at
@ -205,17 +293,15 @@ def snapshot_listing_sources(
)
counters["snapshotted"] = snap_result.rowcount or 0
event_result = db.execute(_EVENT_DIFF_SQL)
counters["price_change_events"] = event_result.rowcount or 0
# Statement возвращает уже готовые пары (counter_key, n) по типам событий —
# dict(...) без Python-агрегации, набор ключей задан инициализацией counters
# выше, так что не сработавшие типы остаются нулями, а не исчезают.
event_rows = db.execute(_EVENT_DIFF_SQL).fetchall()
counters.update(dict(event_rows))
db.commit()
runs_mod.mark_done(db, run_id, counters)
logger.info(
"snapshot_listing_sources run_id=%d done: snapshotted=%d price_change_events=%d",
run_id,
counters["snapshotted"],
counters["price_change_events"],
)
logger.info("snapshot_listing_sources run_id=%d done: %s", run_id, counters)
return counters
except Exception as exc:
logger.exception(

View file

@ -0,0 +1,29 @@
-- 213_listings_snapshots_status_vocab.sql
-- #2674 — словарь listings_snapshots.status стал трёхзначным: 'stale' ≠ 'closed'.
--
-- ПРОБЛЕМА: комментарий колонки (016) обещал два значения — 'active' / 'closed' —
-- и при этом 'closed' не писал никто и никогда: status был константой 'active' у
-- всех 394 704 строк при 55 448 реально неактивных объявлениях. Писатель «снято»
-- появился в #2674, но одним значением обойтись нельзя:
-- - путь 404 (avito_detail_backfill) ЗНАЕТ, что объявления нет: площадка ответила;
-- - путь TTL (deactivate_stale_listings) знает только, что МЫ N суток не смотрели.
-- Замер: прогон по домклику 02.08 деактивировал 6131 объявление за раз (TTL 14 суток
-- против 12 суток простоя обхода) — под общим статусом это 6131 фальшивая «дата
-- продажи» одной датой. Продукт про цены: смешивать факт с догадкой дорого.
--
-- ДЕЛАЕТ: только обновляет COMMENT — сама колонка `text` без CHECK, DDL не нужен.
-- CHECK намеренно НЕ добавляем: 394 704 существующие строки валидны, а жёсткий
-- словарь на историческую таблицу — деструктивный риск ради нулевой выгоды.
--
-- Idempotent: COMMENT ON COLUMN — безусловная перезапись, безопасно повторно.
-- Apply after: 212_sber_index_pull_weekly.sql
BEGIN;
COMMENT ON COLUMN listings_snapshots.status IS
'''active'' = объявление видели в прогоне. '
'''closed'' = площадка ответила 404 на поштучном обходе (жёсткий факт снятия). '
'''stale'' = TTL-деактивация: мы N суток не смотрели (догадка, НЕ дата продажи). '
'NULL = неизвестно. Словарь расширен в #2674 — до него писалось только ''active''.';
COMMIT;

View file

@ -0,0 +1,431 @@
"""#2674 — писатели наконец пишут то, что обещает схема.
Три находки одного класса: колонка есть, писатель есть, тест на писателя зелёный,
а данные не появляются. Обычный юнит-тест такое не ловит по построению он
проверяет то, что автор себе представлял. Ловится это сверкой «что схема обещает»
с «что писатель реально перечисляет», поэтому тесты ниже читают миграции и
сравнивают их с SQL писателя, а не повторяют его же список колонок.
Числа с прода на 2026-08-05/06 (до фикса):
1. house_suggestions 25 055 строк, image_link/area_m2/rooms/floor/total_floors
заполнены у 0 из них (колонки с миграции 064, ~74 дня).
2. listings_snapshots.status 'active' у всех 394 704 строк при 55 448 реально
неактивных объявлений; ни 'closed', ни 'stale' не писал никто и никогда.
3. listing_source_events 8288 строк, все price_change; edited/first_seen
ноль за всё время.
Отдельный класс тестов гейты на то, что писатель НЕ пишет: delisted/relisted схема
разрешает, но при покрытии обхода 10-35% они неотличимы от «скрейпер снова дошёл»
(контроль домклик со 100% покрытием: 0 возвратов за 5 суток), а TTL-путь не имеет
права называть протухание снятием. Журнал и история из догадок хуже пустых.
БД и сеть замоканы реального Postgres не нужно.
"""
from __future__ import annotations
import inspect
import json
import os
import re
from pathlib import Path
from typing import Any
from unittest.mock import AsyncMock, MagicMock, patch
import pytest
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
from scraper_kit.providers.avito.imv import _parse_suggestion
from app.services import house_imv_backfill as hib
from app.tasks import deactivate_stale_avito as deact_mod
from app.tasks import listing_source_snapshot as snap_mod
_SQL_DIR = Path(__file__).resolve().parents[1] / "data" / "sql"
_FIXTURES = Path(__file__).resolve().parent / "fixtures"
# ── Общие хелперы: схема vs writer ────────────────────────────────────────────
def _declared_columns(migration: str, table: str) -> set[str]:
"""Имена колонок из CREATE TABLE IF NOT EXISTS <table> ( ... ); в миграции."""
body = migration.split(f"CREATE TABLE IF NOT EXISTS {table} (")[1].split("\n);")[0]
cols: set[str] = set()
for line in body.splitlines():
m = re.match(r"\s+([a-z_]+)\s+[a-z]", line)
if m:
cols.add(m.group(1))
return cols
def _insert_columns(sql: str, table: str) -> set[str]:
"""Имена колонок из INSERT INTO <table> ( ... ) VALUES."""
m = re.search(rf"INSERT INTO {table}\s*\(([^)]*)\)", sql, re.S)
assert m is not None, f"не найден INSERT INTO {table}"
return {c.strip() for c in m.group(1).split(",") if c.strip()}
# ══ 1. Фотографии подсказок Avito IMV ═════════════════════════════════════════
def test_suggestion_parser_keeps_image_link() -> None:
"""imageLink из ответа площадки доезжает до модели, а не выбрасывается."""
sugg = _parse_suggestion(
{
"id": 8000753763,
"title": "3-к. квартира, 61,6 м², 1/5 эт.",
"price": 7400000,
"imageLink": "https://80.img.avito.st/image/1/abc",
}
)
assert sugg.image_link == "https://80.img.avito.st/image/1/abc"
# В raw_payload ссылку по-прежнему не дублируем — у неё теперь своя колонка.
assert "imageLink" not in (sugg.raw_payload or {})
def test_suggestion_parser_derives_metrics_from_title() -> None:
"""rooms/area_m2/floor/total_floors парсятся из title тем же путём, что у
placementHistory (колонки house_suggestions существуют с миграции 064)."""
sugg = _parse_suggestion(
{"id": 1, "title": "3-к. квартира, 61,6 м², 1/5 эт.", "price": 7400000}
)
assert (sugg.rooms, sugg.area_m2, sugg.floor, sugg.total_floors) == (3, 61.6, 1, 5)
def test_suggestion_parser_survives_unparsable_title() -> None:
"""Нераспознанный заголовок → None'ы, а не исключение (строка всё равно пишется)."""
sugg = _parse_suggestion({"id": 2, "title": "Апартаменты", "price": 1})
assert (sugg.rooms, sugg.area_m2, sugg.floor, sugg.total_floors) == (None, None, None, None)
def test_studio_title_keeps_area_and_floors() -> None:
"""Студия не роняет разбор целиком: 1991 заголовок из 25 055 (7.9%) — без комнатности.
Обязательная группа комнатности обнуляла ВСЕ ЧЕТЫРЕ поля, хотя площадь и этажность
в заголовке есть. rooms=0 конвенция kit'а («0 = студия»), а не «неизвестно».
"""
sugg = _parse_suggestion(
{"id": 3, "title": "Квартира-студия, 34,2 м², 9/10 эт.", "price": 3_500_000}
)
assert (sugg.rooms, sugg.area_m2, sugg.floor, sugg.total_floors) == (0, 34.2, 9, 10)
def test_placement_history_gets_same_title_fix() -> None:
"""Тот же регексп чинит второго писателя — house_placement_history (8.8% без площади)."""
from scraper_kit.providers.avito.imv import _parse_placement_item
item = _parse_placement_item({"id": 9, "title": "Квартира-студия, 28 м², 2/17 эт."})
assert (item.rooms, item.area_m2, item.floor, item.total_floors) == (0, 28.0, 2, 17)
def test_live_fixture_suggestions_carry_image_link() -> None:
"""Живой capture avito_imv_getdata.json: у подсказок реально есть imageLink."""
data = json.loads((_FIXTURES / "avito_imv_getdata.json").read_text("utf-8"))
items = data["suggestions"]["items"]
parsed = [_parse_suggestion(raw) for raw in items]
assert parsed, "фикстура без подсказок — тест бессмыслен"
assert all(s.image_link for s in parsed)
def test_house_suggestions_insert_covers_every_declared_column() -> None:
"""Regression-гейт на весь класс бага: INSERT обязан перечислять КАЖДУЮ колонку
house_suggestions из миграции 064 (кроме автоинкрементного id).
Именно этот тест покраснел бы 74 дня назад: image_link (и заодно area_m2/rooms/
floor/total_floors) объявлены схемой, но в запрос вставки не входили 25 055
строк с NULL. Тест не дублирует список колонок писателя, а сверяет его со схемой,
поэтому ловит и следующую забытую колонку.
"""
declared = _declared_columns(
(_SQL_DIR / "064_house_imv_phase_c.sql").read_text("utf-8"), "house_suggestions"
)
written = _insert_columns(inspect.getsource(hib.save_imv_result), "house_suggestions")
assert (
declared - {"id"} <= written
), f"колонки без писателя: {sorted(declared - {'id'} - written)}"
def test_save_imv_result_binds_image_link_and_metrics() -> None:
"""save_imv_result передаёт значения подсказки в bind-параметры (не только в SQL)."""
sugg = _parse_suggestion(
{
"id": 777,
"title": "2-к. квартира, 42 м², 4/5 эт.",
"price": 6300000,
"imageLink": "https://img/x.jpg",
}
)
result = MagicMock(
cache_key="k",
recommended_price=1,
lower_price=1,
higher_price=1,
market_count=1,
raw_response=None,
placement_history=[],
suggestions=[sugg],
)
params = {
"rooms": 2,
"area_m2": 42.0,
"floor": 4,
"floor_at_home": 5,
"house_type": "panel",
"renovation_type": "cosmetic",
"has_balcony": True,
"has_loggia": False,
}
db = MagicMock()
hib.save_imv_result(db, house_id=1, params=params, result=result)
sugg_calls = [
c for c in db.execute.call_args_list if "INSERT INTO house_suggestions" in str(c.args[0])
]
assert len(sugg_calls) == 1
bound = sugg_calls[0].args[1]
assert bound["img"] == "https://img/x.jpg"
assert (bound["rooms"], bound["area"], bound["floor"], bound["total_floors"]) == (2, 42.0, 4, 5)
# ══ 2. «Снято» и «протухло» в дневной истории объявлений ══════════════════════
_DEACT_SQL_BUILDERS = (deact_mod._build_all_segments_sql, deact_mod._build_segments_sql)
@pytest.mark.parametrize("build", _DEACT_SQL_BUILDERS)
@pytest.mark.parametrize("column", sorted(deact_mod._ALLOWED_STALENESS_COLUMNS))
def test_deactivation_writes_stale_snapshot_in_same_statement(build: Any, column: str) -> None:
"""Деактивация и снимок 'stale' — один statement, значит одна транзакция.
До #2674 задача только двигала флаг: в listings_snapshots не появлялось ничего,
и дата снятия объявления (лучший сигнал «скорее всего продано») восстанавливалась
на глаз из последнего показа + предполагаемого срока жизни.
"""
sql = str(build(column).text)
assert "SET is_active = false" in sql
assert "RETURNING id, price_rub" in sql
assert "INSERT INTO listings_snapshots" in sql
assert "'stale'" in sql
# Снимок пишется по строкам, которые вернул сам UPDATE, — не отдельной выборкой.
assert "FROM stale" in sql
# Идемпотентность: повторный прогон в те же сутки не падает на PK.
assert "ON CONFLICT (listing_id, snapshot_date) DO UPDATE" in sql
# psycopg v3: никаких :param::type.
assert not re.search(r":\w+::", sql)
@pytest.mark.parametrize("build", _DEACT_SQL_BUILDERS)
def test_ttl_path_never_claims_closed(build: Any) -> None:
"""TTL-путь НЕ имеет права писать 'closed' — он не знает, что объявление снято.
Замер: прогон по домклику 02.08 деактивировал 6131 объявление за раз (TTL 14 суток
против 12 суток простоя обхода). Под статусом 'closed' это 6131 фальшивая «дата
продажи» одной датой. 'closed' остаётся только за 404 там ответила площадка.
"""
assert "'closed'" not in str(build("last_seen_at").text)
class _FakeResult:
def __init__(self, rowcount: int) -> None:
self.rowcount = rowcount
class _FakeDB:
def __init__(self, rowcount: int = 0) -> None:
self._rowcount = rowcount
self.executed: list[tuple[Any, Any]] = []
self.committed = False
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
self.executed.append((stmt, params))
return _FakeResult(self._rowcount)
def commit(self) -> None:
self.committed = True
def rollback(self) -> None: # pragma: no cover — путь ошибки тут не проверяется
pass
def test_deactivate_stale_listings_threads_run_id_into_snapshot(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""run_id доезжает до снимка — провенанс «каким прогоном закрыто» не теряется."""
monkeypatch.setattr(deact_mod.runs_mod, "mark_done", lambda *a, **k: None)
monkeypatch.setattr(deact_mod.runs_mod, "mark_failed", lambda *a, **k: None)
db = _FakeDB(rowcount=7)
out = deact_mod.deactivate_stale_listings(
db, # type: ignore[arg-type]
run_id=555,
listing_source="avito",
ttl_days=10,
)
assert out == {"deactivated": 7}
assert len(db.executed) == 1, "деактивация и снимок обязаны быть одним statement'ом"
stmt, bound = db.executed[0]
assert "INSERT INTO listings_snapshots" in str(stmt)
assert bound is not None and bound["run_id"] == 555
@pytest.mark.asyncio
async def test_avito_404_records_closed_snapshot() -> None:
"""404 с площадки — самый достоверный сигнал снятия; он тоже попадает в историю."""
from scraper_kit.avito_exceptions import AvitoListingGoneError
from app.tasks.avito_detail_backfill import run_avito_detail_backfill
db = MagicMock()
sel = MagicMock()
sel.mappings.return_value.all.return_value = [
{"id": 42, "source_url": "/items/42", "price_rub": 5_000_000}
]
db.execute.return_value = sel
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False)
with (
patch("app.tasks.avito_detail_backfill.settings", fake_settings),
patch("app.tasks.avito_detail_backfill.AsyncSession", return_value=AsyncMock()),
patch("app.tasks.avito_detail_backfill.AvitoScraper"),
patch("app.tasks.avito_detail_backfill.runs_mod", MagicMock()),
patch("app.tasks.avito_detail_backfill.asyncio.sleep", new_callable=AsyncMock),
patch(
"app.tasks.avito_detail_backfill.fetch_detail",
AsyncMock(side_effect=AvitoListingGoneError("404 gone")),
),
patch("app.tasks.avito_detail_backfill.upsert_listing_snapshot") as snap,
):
result = await run_avito_detail_backfill(db, run_id=3, params={"budget_sec": 60})
assert result.gone == 1
snap.assert_called_once()
assert snap.call_args.kwargs["status"] == "closed"
assert snap.call_args.kwargs["listing_id"] == 42
assert snap.call_args.kwargs["price_rub"] == 5_000_000
# ══ 3. Журнал событий объявлений ══════════════════════════════════════════════
def _schema_event_types() -> set[str]:
"""Пять типов из CHECK-констрейнта миграции 079 — источник правды."""
sql = (_SQL_DIR / "079_listing_source_history.sql").read_text("utf-8")
check = sql.split("event_type IN (")[1].split(")")[0]
return set(re.findall(r"'([a-z_]+)'", check))
# Два типа схемы НЕ ВЫВОДИМЫ из наших данных и намеренно не пишутся (#2674).
# is_active в снимке значит «мы видели», а не «есть на площадке», поэтому переход
# рождается тем, что скрейпер снова дошёл до источника. Контрольная группа за 14-18.07:
# domklik при покрытии 99.9-100% дал возвратов РОВНО 0 и снятий 1-4 в сутки, yandex при
# 34-43% — снятий 343-433 в сутки. Тот же обход, тот же день, разница только в покрытии.
# Отсюда: avito 13.07 (остановка обхода) 3023 «снятия» за сутки против контрольных 1-4
# (точность ≈4%), и 4705 «возвратов» из 5493 за 12 дней — два дня после возобновления.
_NOT_DERIVABLE_EVENT_TYPES = {"delisted", "relisted"}
def test_event_writer_covers_every_derivable_schema_event_type() -> None:
"""Писатель обязан уметь каждый ВЫВОДИМЫЙ тип из CHECK схемы.
До #2674 из пяти типов писался один (price_change, 8288 строк). Дописаны два
выводимых; два оставшихся сознательное решение, а не забытая ветка (см.
_NOT_DERIVABLE_EVENT_TYPES). Тест сверяет со схемой, а не с копией списка,
поэтому покраснеет и на шестом типе, добавленном в CHECK без писателя.
"""
declared = _schema_event_types()
assert len(declared) == 5, f"схема 079 изменилась: {sorted(declared)}"
sql = str(snap_mod._EVENT_DIFF_SQL.text)
expected = declared - _NOT_DERIVABLE_EVENT_TYPES
missing = {t for t in expected if f"'{t}'" not in sql}
assert not missing, f"выводимые типы без писателя: {sorted(missing)}"
def test_not_derivable_events_are_never_written() -> None:
"""delisted/relisted не пишутся: при покрытии обхода 10-35% они неотличимы от
«скрейпер снова дошёл». Контроль домклик со 100% покрытием: 0 возвратов за 5 суток.
Гейт против «дописать для полноты»: журнал из догадок хуже пустого журнала.
"""
sql = str(snap_mod._EVENT_DIFF_SQL.text)
written = {t for t in _NOT_DERIVABLE_EVENT_TYPES if f"'{t}'" in sql}
assert not written, f"невыводимые типы попали в писатель: {sorted(written)}"
# is_active больше не читается вовсе — иначе ветка вернётся незаметно.
assert "is_active" not in sql
def test_first_seen_requires_left_join_and_derivations_use_snapshot_fields() -> None:
"""Ветки выводятся из полей снимка, first_seen достижим только через LEFT JOIN.
С обычным JOIN источник без предыдущего снимка отбрасывался джойном событие
«первое появление» было недостижимо по построению.
"""
sql = str(snap_mod._EVENT_DIFF_SQL.text)
assert "LEFT JOIN LATERAL" in sql
assert "p.snapshot_date IS NULL" in sql # first_seen
assert "t.payload_hash IS DISTINCT FROM p.payload_hash" in sql # edited
# Снимок за сегодня обязан отдавать поля, из которых выводятся ветки.
assert "SELECT listing_source_id, price_rub, payload_hash" in sql
def test_event_dedup_works_across_runs_not_only_within_one() -> None:
"""change_time усечён до суток: UNIQUE(source, change_time, type) должен гасить
повторный прогон в те же сутки (2 августа их было два), а не только строки одного."""
sql = str(snap_mod._EVENT_DIFF_SQL.text)
assert "date_trunc('day', now())" in sql
def test_price_change_division_guarded_by_nullif() -> None:
"""VALUES вычисляется ДО фильтра e.fires → без NULLIF прогон падал бы на
первом источнике с нулевой прошлой ценой (предикат p.price_rub <> 0 не спасает)."""
sql = str(snap_mod._EVENT_DIFF_SQL.text)
assert "NULLIF(p.price_rub, 0)" in sql
class _EventFakeDB:
"""Session-заглушка: SET LOCAL и snapshot дают rowcount, event-diff — пары счётчиков."""
def __init__(self, event_rows: list[tuple[str, int]]) -> None:
self._event_rows = event_rows
self.committed = False
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
result = MagicMock()
if "listing_source_events" in str(stmt):
result.fetchall.return_value = self._event_rows
else:
result.rowcount = 100
return result
def commit(self) -> None:
self.committed = True
def rollback(self) -> None: # pragma: no cover — путь ошибки тут не проверяется
pass
def test_counters_report_every_written_type_including_zeros(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""Счётчики прогона показывают все пишущиеся типы; не сработавший честно равен 0.
Ровно этого счётчика не хватало, чтобы заметить четыре нуля из пяти за 66 дней.
Счётчиков НЕвыводимых типов быть не должно иначе вечный 0 будет читаться как
«событий не было», а не как «мы это сознательно не пишем».
"""
monkeypatch.setattr(snap_mod.runs_mod, "mark_done", lambda *a, **k: None)
monkeypatch.setattr(snap_mod.runs_mod, "mark_failed", lambda *a, **k: None)
db = _EventFakeDB([("first_seen_events", 600), ("edited_events", 17)])
out = snap_mod.snapshot_listing_sources(db, run_id=1) # type: ignore[arg-type]
assert out["first_seen_events"] == 600
assert out["edited_events"] == 17
assert out["price_change_events"] == 0
# Ключ на каждый пишущийся тип — иначе «ноль» неотличим от «типа нет в counters».
for event_type in _schema_event_types() - _NOT_DERIVABLE_EVENT_TYPES:
assert f"{event_type}_events" in out
for event_type in _NOT_DERIVABLE_EVENT_TYPES:
assert f"{event_type}_events" not in out

View file

@ -104,9 +104,11 @@ def test_event_diff_is_set_based_lateral_not_python_loop() -> None:
def test_event_diff_emits_price_change_with_diff_percent() -> None:
assert "'price_change'" in _EVENT_DIFF_SQL
# diff_percent = (new-old)/old*100.
# diff_percent = (new-old)/old*100. NULLIF на знаменателе (#2674): выражения
# VALUES вычисляются ДО фильтра e.fires, поэтому предикат "p.price_rub <> 0"
# больше не защищает само деление.
assert "(t.price_rub - p.price_rub)" in _EVENT_DIFF_SQL
assert "/ p.price_rub * 100" in _EVENT_DIFF_SQL
assert "/ NULLIF(p.price_rub, 0) * 100" in _EVENT_DIFF_SQL
# Only when the price actually changed and old is a usable denominator.
assert "t.price_rub <> p.price_rub" in _EVENT_DIFF_SQL
assert "p.price_rub <> 0" in _EVENT_DIFF_SQL
@ -200,22 +202,30 @@ def test_migration_079_uses_psycopg_safe_sql() -> None:
class _FakeResult:
def __init__(self, rowcount: int) -> None:
def __init__(self, rowcount: int, rows: list[tuple[str, int]] | None = None) -> None:
self.rowcount = rowcount
self._rows = rows or []
def fetchall(self) -> list[tuple[str, int]]:
"""Event-diff statement возвращает пары (counter_key, n) — см. #2674."""
return self._rows
class _FakeDB:
"""Minimal stand-in for a SQLAlchemy Session — records execute() calls, returns rowcounts."""
def __init__(self, rowcounts: list[int]) -> None:
def __init__(
self, rowcounts: list[int], event_rows: list[tuple[str, int]] | None = None
) -> None:
self._rowcounts = list(rowcounts)
self._event_rows = event_rows or []
self.executed: list[Any] = []
self.committed = False
self.rolled_back = False
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
self.executed.append((stmt, params))
return _FakeResult(self._rowcounts.pop(0))
return _FakeResult(self._rowcounts.pop(0), self._event_rows)
def commit(self) -> None:
self.committed = True
@ -225,7 +235,12 @@ class _FakeDB:
def test_counter_logic_with_fake_db(monkeypatch: pytest.MonkeyPatch) -> None:
"""snapshot_listing_sources maps the two execute() rowcounts to its counters + marks done."""
"""snapshot_listing_sources maps the snapshot rowcount + per-type event rows to counters.
#2674: event-diff больше не отдаёт один rowcount — внешний SELECT над
data-modifying CTE возвращает пары (counter_key, n) по типам событий, а
counters инициализированы всеми пятью ключами (не сработавший тип = честный 0).
"""
marked: dict[str, Any] = {}
monkeypatch.setattr(
snap_mod.runs_mod,
@ -234,11 +249,16 @@ def test_counter_logic_with_fake_db(monkeypatch: pytest.MonkeyPatch) -> None:
)
monkeypatch.setattr(snap_mod.runs_mod, "mark_failed", lambda *a, **k: None)
# rowcounts: SET LOCAL statement_timeout (ignored), snapshot upsert, event-diff insert.
db = _FakeDB(rowcounts=[0, 18355, 42])
# rowcounts: SET LOCAL statement_timeout (ignored), snapshot upsert, event-diff.
db = _FakeDB(rowcounts=[0, 18355, 0], event_rows=[("price_change_events", 42)])
out = snap_mod.snapshot_listing_sources(db, run_id=99) # type: ignore[arg-type]
assert out == {"snapshotted": 18355, "price_change_events": 42}
assert out == {
"snapshotted": 18355,
"price_change_events": 42,
"edited_events": 0,
"first_seen_events": 0,
}
assert db.committed is True
assert len(db.executed) == 3
# First statement sets the per-transaction wall-clock budget (#2607).
@ -249,7 +269,7 @@ def test_counter_logic_with_fake_db(monkeypatch: pytest.MonkeyPatch) -> None:
assert params is not None and params["run_id"] == 99
# Run finalised via mark_done with the same counters.
assert marked["run_id"] == 99
assert marked["counters"] == {"snapshotted": 18355, "price_change_events": 42}
assert marked["counters"] == out
# ── budget_sec / statement_timeout (#2607) ─────────────────────────────────────

View file

@ -141,9 +141,18 @@ class IMVSuggestion:
title: str | None = None
address: str | None = None
price_rub: int | None = None
# rooms/area_m2/floor/total_floors парсятся из title (#2674) — тот же путь, что
# у IMVPlacementHistoryItem; колонки в house_suggestions есть с миграции 064.
rooms: int | None = None
area_m2: float | None = None
floor: int | None = None
total_floors: int | None = None
exposure_days: int | None = None
publish_date: date | None = None
item_url: str | None = None
# image_link — ссылка на фото лота (raw imageLink). Колонка house_suggestions.image_link
# существует с 064, но до #2674 значение выбрасывалось парсером.
image_link: str | None = None
metro_name: str | None = None
metro_distance: str | None = None
metro_color: str | None = None
@ -206,10 +215,18 @@ def compute_imv_cache_key(
# Паттерн для заголовка вида "2-к. квартира, 42 м², 4/5 эт."
# Комнатность НЕОБЯЗАТЕЛЬНА (#2674): 1991 заголовок из 25 055 (7.9%) — «Квартира-студия,
# 34,2 м², 9/10 эт.». Площадь и этажность там есть, но обязательная группа комнатности
# роняла match целиком и обнуляла ВСЕ ЧЕТЫРЕ поля. Тот же потолок был виден на соседней
# таблице (8.8% строк house_placement_history без площади) — обоих писателей чинит один
# регексп. Опциональная группа жадная, поэтому «3-к. квартира…» по-прежнему даёт rooms=3.
_TITLE_RE = re.compile(
r"^(?P<rooms>\d+)-к[.\s].*?(?P<area>[\d,]+)\s*м².*?(?P<floor>\d+)/(?P<total>\d+)\s*эт",
r"^(?:(?P<rooms>\d+)-к[.\s])?.*?(?P<area>[\d,]+)\s*м².*?(?P<floor>\d+)/(?P<total>\d+)\s*эт",
re.IGNORECASE,
)
# Студия = 0 комнат — конвенция kit'а (scraper_kit.base.RawLot.rooms «0 = студия»,
# providers/yandex/detail.py). Отличаем её от «комнатность неизвестна» (None).
_TITLE_STUDIO_RE = re.compile(r"студи[яюей]", re.IGNORECASE)
def _parse_title(title: str | None) -> dict[str, int | float | None]:
@ -226,7 +243,11 @@ def _parse_title(title: str | None) -> dict[str, int | float | None]:
if not m:
return result
try:
result["rooms"] = int(m.group("rooms"))
rooms = m.group("rooms")
if rooms:
result["rooms"] = int(rooms)
elif _TITLE_STUDIO_RE.search(title):
result["rooms"] = 0 # студия
result["area_m2"] = float(m.group("area").replace(",", "."))
result["floor"] = int(m.group("floor"))
result["total_floors"] = int(m.group("total"))
@ -262,20 +283,28 @@ def _parse_placement_item(raw: dict[str, Any]) -> IMVPlacementHistoryItem:
def _parse_suggestion(raw: dict[str, Any]) -> IMVSuggestion:
"""Парсит один элемент из suggestions.items."""
title = raw.get("title")
parsed = _parse_title(title)
metro = raw.get("metro") or {}
colors: list[str] = metro.get("colors") or []
# Фильтруем imageLink из raw_payload
# imageLink выносим в отдельное поле (колонка house_suggestions.image_link) и
# убираем из raw_payload, чтобы не хранить ссылку дважды.
clean_raw = {k: v for k, v in raw.items() if k not in ("imageLink",)}
return IMVSuggestion(
ext_item_id=str(raw["id"]),
title=raw.get("title"),
title=title,
address=raw.get("address"),
price_rub=raw.get("price"),
rooms=parsed["rooms"], # type: ignore[arg-type]
area_m2=parsed["area_m2"], # type: ignore[arg-type]
floor=parsed["floor"], # type: ignore[arg-type]
total_floors=parsed["total_floors"], # type: ignore[arg-type]
exposure_days=raw.get("exposure"),
publish_date=_unix_to_date(raw.get("publishDate")),
item_url=raw.get("itemLink"),
image_link=raw.get("imageLink"),
metro_name=metro.get("name"),
metro_distance=metro.get("distance"),
metro_color=colors[0] if colors else None,

View file

@ -4,11 +4,16 @@
PRIMARY KEY (listing_id, snapshot_date) максимум 1 snapshot в сутки.
ON CONFLICT DO UPDATE берём последний за день (перезаписываем при повторном run'е).
Используется двумя путями:
Используется тремя путями:
1. SERP scrape (save_listings в base.py) записывает цену + позицию в выдаче
каждый раз когда listing появляется в поиске.
каждый раз когда listing появляется в поиске (status='active').
2. Detail backfill (save_detail_enrichment в cian_detail.py) записывает цену
из detail-страницы для listings которые раньше не имели snapshot'а.
из detail-страницы для listings которые раньше не имели snapshot'а (status='active').
3. Снятие объявления (#2674): avito_detail_backfill при 404 с площадки пишет
status='closed'. Второй писатель 'closed' deactivate_stale_listings, он
набирает тысячи строк за прогон и пишет их одним set-based statement'ом
(data-modifying CTE в app/tasks/deactivate_stale_avito.py), не через этот
per-row хелпер.
"""
from __future__ import annotations