fix(tradein): выборка миграции 286 повторяет гейт 1:1, правило первой точки (#3376)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 7s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 10s
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) Successful in 4m56s

Ревью нашло, что миграция и код ловили РАЗНОЕ. Миграция брала базой предыдущую
СЫРУЮ строку (lag), гейт — предыдущую ОСТАВЛЕННУЮ. На 1M→10M→1M→10M (цена 1M)
lag-версия удаляла честную точку, на 1M→10M→1.05M→9.9M — не была идемпотентной
(второй прогон доедал 9.9M). Теперь кандидаты выбирает PL/pgSQL-цикл, пошагово
повторяющий drop_decimal_slips, а правило первой точки — отдельным INSERT..SELECT
уже по ОСТАВШИМСЯ строкам.

Правило первой точки — из прод-разбора: 12 из 20 остатков domklik это серии вида
330 000 → 3 300 000 (текущая цена 3 300 000) и 420 000 → 4 200 000 → 4 500 000,
где дефектная точка ПЕРВАЯ и базы слева у неё нет. Свидетелей по-прежнему два:
×10 ко второй точке И подтверждение второй третьей-или-текущей-ценой. Решение по
первой точке принимается по kept-серии, а не по сырой, — иначе гейт теряет
идемпотентность (перебор ловит 1122 таких прогона).

Идемпотентность доказана НА ГЕЙТЕ: property-тест gate(gate(s)) == gate(s) по всем
сериям длины 2-6 (19 525 серий × 3 текущие цены). Фальсифицирован обеими
поломками — сырая база даёт 136 красных прогонов, сырые соседи первой точки 1122.
Раз SQL зеркалит гейт, свойство переносится на миграцию.

Ещё в 286: третий свидетель ПРОТИВ удаления (цена подтверждена триггерной строкой
того же объявления — значит она реально наблюдалась в listings.price_rub) и
финальный шаг |diff_percent| > 100 → NULL по всем источникам, то же правило, что
validate_diff_percent на записи. Ожидаемое число удалений в шапке — 35 + ~12 из
двухсвидетельского предзамера, а не 76 (то была односвидетельская цифра).
Прогон обеих фаз дважды с ROLLBACK — tradein-mvp/scripts/sql/286_dryrun.sql.

yandex: проводка гейта снята как мёртвая. На том пути серия из двух точек, а
свидетель последней — текущая цена лота, то есть она же сама: ветка по построению
не могла выбросить ничего. Оставлен честный комментарий-потолок и ссылка на
follow-up (отлов требует DELETE на следующем наблюдении).

cian: после выброса точки соседу пересчитывается diff_percent (было только у
domclick). domclick: цена листинга берётся RETURNING'ом у UPDATE вместо отдельного
SELECT по PK, пересчёт вынесен в общий recompute_diff_percent с гейтом на пустую
цену (ручной ingest кладёт price_changes из JSONL без валидации).
This commit is contained in:
bot-backend 2026-09-06 03:53:08 +05:00
parent 142967d064
commit 721ceb9876
9 changed files with 810 additions and 127 deletions

View file

@ -34,7 +34,6 @@ from datetime import UTC, datetime, timedelta
# isinstance checks — так что любой ScrapedLot-совместимый объект (source, source_id,
# price_rub, price_previous_rub, price_trend) остаётся взаимозаменяем для этой функции.
from scraper_kit.base import ScrapedLot
from scraper_kit.offer_price_history import drop_decimal_slips
from sqlalchemy import text
from sqlalchemy.orm import Session
@ -117,7 +116,6 @@ def record_yandex_price_history(db: Session, lots: list[ScrapedLot]) -> int:
inserted = 0
skipped = 0
errors = 0
slips = 0
for lot in lots:
if lot.source != _SOURCE:
@ -142,47 +140,46 @@ def record_yandex_price_history(db: Session, lots: list[ScrapedLot]) -> int:
{"listing_id": listing_id},
).fetchone()
# Серия для гейта сдвигов разряда (#3376): (цена, change_time).
# change_time=None — якорь, последняя лежащая в БД точка: гейту нужна
# как база отсчёта, вставлять её повторно не надо.
points: list[tuple[float, datetime | None]] = []
# ПОТОЛОК: гейта сдвига разряда (#3376) здесь НЕТ, и он бы тут не
# сработал. drop_decimal_slips требует двух свидетелей — скачка ×10
# к предыдущей точке и подтверждения у следующей. На этом пути серия
# максимум из двух точек (последняя лежащая в БД + текущая), у
# последней точки свидетель — текущая цена лота, а она и ЕСТЬ эта
# точка: свидетель совпадает с подозреваемым, отношение всегда 1.0.
# То есть проводка была бы декорацией: ветка, которая по построению
# не может выбросить ни одной точки. Отлов ×10 у yandex требует
# другого механизма — сравнения со СЛЕДУЮЩИМ наблюдением, а значит
# DELETE уже вставленной строки. Отдельная задача: #TODO-yandex-slip.
if latest is None:
# История пуста — seed (+ опционально previous-точка).
prev = lot.price_previous_rub
if prev is not None and prev != lot.price_rub:
points.append((float(prev), prev_ts))
points.append((float(lot.price_rub), now))
else:
latest_price = float(latest[0])
if latest_price == float(lot.price_rub):
skipped += 1
continue
points.append((latest_price, None))
points.append((float(lot.price_rub), now))
# Потолок гейта на этом пути: у yandex последняя точка серии И есть
# текущая цена лота, то есть второй свидетель совпадает с самой
# точкой и она не выбрасывается никогда. Ловится только сдвиг ВНУТРИ
# серии (seed из price.previous + current). Уже лежащие в БД
# yandex-сдвиги чистит миграция 286, живой отлов «на следующем
# наблюдении» потребовал бы DELETE — отдельная задача.
kept, dropped = drop_decimal_slips(
points,
lambda point: point[0],
current_price=lot.price_rub,
listing_id=listing_id,
)
slips += dropped
for price_rub, change_time in kept:
if change_time is None:
continue
_insert_point(
db,
listing_id=listing_id,
price_rub=prev,
change_time=prev_ts,
)
inserted += 1
_insert_point(
db,
listing_id=listing_id,
price_rub=price_rub,
change_time=change_time,
price_rub=lot.price_rub,
change_time=now,
)
inserted += 1
else:
latest_price = float(latest[0])
if latest_price != float(lot.price_rub):
_insert_point(
db,
listing_id=listing_id,
price_rub=lot.price_rub,
change_time=now,
)
inserted += 1
else:
skipped += 1
except Exception as exc:
# Per-lot SAVEPOINT откатился — продолжаем батч, не роняем остальные lot'ы.
errors += 1
@ -194,11 +191,10 @@ def record_yandex_price_history(db: Session, lots: list[ScrapedLot]) -> int:
db.commit()
logger.info(
"yandex_price_history: lots=%d inserted=%d skipped=%d errors=%d decimal_slips_dropped=%d",
"yandex_price_history: lots=%d inserted=%d skipped=%d errors=%d",
len(lots),
inserted,
skipped,
errors,
slips,
)
return inserted

View file

@ -1,33 +1,55 @@
-- 286_offer_price_history_decimal_slips.sql
-- Удалить из offer_price_history точки с потерянным разрядом и починить соседа (#3376).
-- Удалить из offer_price_history точки с потерянным разрядом и починить diff_percent (#3376).
--
-- ЧТО НЕ ТАК
-- В priceHistory источников встречаются точки ровно ×10 / ÷10 к соседям:
-- 377 000 → 3 770 000 → 3 720 000; 19 064 000 → 190 150 000 → 19 150 000. Цена
-- возвращается к базе следующей же точкой — это потерянный разряд у источника,
-- а не рынок. Замер на проде 06.09.2026 (LAG по listing_id, отношение к
-- предыдущей точке в [9.8, 10.2] либо [0.098, 0.102]):
-- yandex 30 961 строк → 6 вверх, 0 вниз, 6 объявлений (все — строки загрузчика)
-- domklik 22 020 строк → 30 вверх, 18 вниз, 35 объявлений (все — строки загрузчика)
-- cian 18 266 строк → 9 вверх, 12 вниз, 13 объявлений (все — строки загрузчика)
-- avito 6 406 строк → 0 вверх, 1 вниз, 1 объявление (строка ТРИГГЕРА)
-- а не рынок. Второй вид того же дефекта — в НАЧАЛЕ серии: 330 000 → 3 300 000
-- (текущая цена объявления 3 300 000), 420 000 → 4 200 000 → 4 500 000. Дефектная
-- точка первая, слева у неё базы нет, поэтому её ловит зеркальное правило.
--
-- Читатели колонки — медианный торг лендинга (#3223), админка, /scrapers,
-- публичный API. Оценщик (services/estimator.py) историю не читает, на оценку
-- эти точки не влияли.
--
-- Код починен в том же PR: гейт drop_decimal_slips в scraper_kit/offer_price_history.py
-- подключён ко всем писателям истории (domclick/detail.py, cian/detail.py,
-- backend/app/services/yandex_price_history.py). Эта миграция отрабатывает
-- задним числом по уже собранным строкам.
-- подключён к писателям истории (domclick/detail.py, cian/detail.py). Эта миграция
-- отрабатывает задним числом по уже собранным строкам.
--
-- КРИТЕРИЙ — ТОТ ЖЕ, ЧТО В КОДЕ, И ТРЕБУЕТ ДВУХ СВИДЕТЕЛЕЙ
-- 1) скачок к предыдущей точке: price/prev ∈ [9.5, 10.5] ИЛИ prev/price ∈ [9.5, 10.5];
-- 2) возврат к базе у следующей точки: next/prev ∈ [0.9, 1.1].
-- ВЫБОРКА КАНДИДАТОВ ЗЕРКАЛИТ ГЕЙТ 1:1, А НЕ «ТОТ ЖЕ КРИТЕРИЙ»
-- Первая редакция этого файла брала предыдущую точку через lag(), то есть СЫРУЮ
-- соседнюю строку — включая те, что удаляются этим же проходом. Гейт в коде берёт
-- предыдущую ОСТАВЛЕННУЮ. Расхождение видно на двух сериях (обе — от ревьюера):
-- 1M → 10M → 1M → 10M (текущая цена 1M): гейт оставляет [1M, 1M], lag-версия
-- оставляла [1M] — честная точка удалена;
-- 1M → 10M → 1.05M → 9.9M (текущая цена 1M): lag-версия удаляла на первом проходе
-- ещё и 1.05M (её сырая база — уже удалённая 10M), а на втором доедала 9.9M,
-- то есть НЕ была идемпотентной.
-- Поэтому кандидаты выбираются PL/pgSQL-циклом, который повторяет гейт пошагово:
-- строки по (change_time, id), переменная base — последняя ОСТАВЛЕННАЯ цена,
-- свидетель — следующая СЫРАЯ точка либо, у последней, listings.price_rub.
-- Правило первой точки — вторым, УЖЕ ПО ОСТАВШИМСЯ строкам (в гейте оно тоже
-- применяется к kept-серии: свидетели первой точки сами могут оказаться сдвигами).
--
-- КРИТЕРИЙ — ДВА СВИДЕТЕЛЯ, И У ПЕРВОЙ ТОЧКИ ТОЖЕ ДВА
-- Внутри серии: 1) скачок к предыдущей ОСТАВЛЕННОЙ точке, price/base ∈ [9.5, 10.5]
-- в любую сторону; 2) возврат к базе у следующей точки, next/base ∈ [0.9, 1.1].
-- Первая точка: 1) она ×10/÷10 ко второй оставшейся; 2) вторая подтверждена третьей
-- оставшейся, а если третьей нет — текущей ценой объявления, witness/second ∈ [0.9, 1.1].
-- Окно шире чистой десятки, потому что сдвиг разряда часто идёт вместе с настоящим
-- мелким изменением цены (377 000 → 3 720 000 это ×9.87). Одного скачка НЕ хватает:
-- ×10 без возврата к базе — это возможное честное изменение цены, и такие точки
-- остаются. У последней точки листинга следующей нет — свидетелем берём текущую
-- цену объявления (listings.price_rub); нет и её → точку не трогаем.
-- ×10 без подтверждения — это возможное честное изменение цены (400 000 → 4 000 000 →
-- 8 000 000 — разгон, а не разряд), и такие точки остаются. Нет второго свидетеля
-- вообще (последняя точка серии, а listings.price_rub пуст) → точку не трогаем.
--
-- СКОЛЬКО ЖДЁМ УДАЛЕНИЙ
-- Предзамер на проде 06.09.2026 под ДВУХСВИДЕТЕЛЬСКИМ критерием: 35 строк внутри
-- серий + ~12 первых точек (12 из 20 остатков domklik — ровно вид 330 000 → 3 300 000)
-- ≈ 47. Прежняя цифра 76 в шапке была из ОДНОСВИДЕТЕЛЬСКОГО замера (просто отношение
-- к предыдущей строке) и к этому критерию отношения не имеет. Порог остановки 200 —
-- четырёхкратный запас: поймали больше — критерий ловит не то, и миграция обязана
-- упасть, а не молча вычистить историю.
--
-- ТОЛЬКО СТРОКИ ЗАГРУЗЧИКА: change_time <> recorded_at
-- Признак происхождения тот же, что в 285 (там же его доказательство): триггер
@ -35,65 +57,169 @@
-- INSERT'ом → метки равны побайтово; загрузчик кладёт в change_time дату источника,
-- а recorded_at оставляет DEFAULT NOW(). Триггерные строки (миграция 131) — это
-- зафиксированные ЖИВЫЕ смены listings.price_rub, а не разбор чужого массива; их
-- не удаляем и в окно lag()/lead() не берём. Поэтому единственная avito-строка из
-- замера остаётся: она триггерная, и если цена в карточке действительно менялась
-- ×10, то это наблюдение, а не дефект разбора.
-- не удаляем и в цепочку не берём. Поэтому единственная avito-строка из замера
-- остаётся: она триггерная, и если цена в карточке действительно менялась ×10, то
-- это наблюдение, а не дефект разбора.
--
-- ОТСЮДА ЖЕ ТРЕТИЙ СВИДЕТЕЛЬ, КОТОРЫЙ СПАСАЕТ ОТ УДАЛЕНИЯ (L-4): если ровно такая
-- же цена есть у ТРИГГЕРНОЙ строки того же объявления — значит listings.price_rub
-- когда-то реально равнялась этому значению, и подозреваемая точка загрузчика
-- подтверждена независимым писателем. Такую не удаляем и оставляем в цепочке базой.
--
-- ПЕРЕСЧЁТ diff_percent У СЛЕДУЮЩЕЙ СТРОКИ
-- У строки, шедшей за удалённой, предыдущая цена сменилась — процент, посчитанный
-- от фантомной базы, надо пересчитать по формуле 285: (price prev)/prev*100,
-- prev отсутствует/нулевой → NULL (процента не существует, ноль читался бы как
-- «цена не менялась»). |x| > 100 → NULL — ровно то, что делает validate_diff_percent
-- на записи: такое значение уже не процент.
-- «цена не менялась»). Цепочки подряд удалённых строк это покрывает: next_id
-- последней удалённой в цепочке и есть первая уцелевшая.
--
-- ФИНАЛЬНЫЙ ШАГ: |diff_percent| > 100 → NULL ПО ВСЕМ ИСТОЧНИКАМ
-- Ровно то, что делает validate_diff_percent на записи с #3225: такая величина —
-- уже не процент (у domklik туда клали рубли). 285 пересчитала только domklik и
-- только из соседних цен, поэтому у cian/yandex/avito-строк загрузчика мусор мог
-- остаться. Честные скачки ×2…×6 при этом СОХРАНЯЮТ строку и цену — обнуляется
-- только процент, посчитанный не от той величины.
--
-- ИДЕМПОТЕНТНОСТЬ
-- Второй прогон: удалённых точек больше нет → выборка пуста → 0 удалений и 0
-- пересчётов (UPDATE вдобавок ограничен IS DISTINCT FROM). Новых объектов схемы нет,
-- временная таблица уходит по ON COMMIT DROP.
--
-- СТРАХОВКА
-- Замер дал 76 строк-кандидатов. Если критерий поймал больше 200 — он ловит не то
-- (или данные разъехались), и миграция обязана упасть, а не молча вычистить историю.
-- Свойство доказано НА ГЕЙТЕ, а не на этом файле: property-тест
-- test_drop_decimal_slips_is_idempotent перебирает серии длины 2-6 по алфавиту
-- {1M, 10M, 100M, 1.05M, 9.9M} × current_price ∈ {—, 1M, 10M} и требует
-- gate(gate(s)) == gate(s). Раз выборка кандидатов здесь повторяет гейт пошагово,
-- свойство переносится: второй прогон даёт 0 удалений. Пересчёт diff_percent
-- детерминирован и ограничен IS DISTINCT FROM, финальный NULL-шаг — тоже (после
-- него |x| > 100 не остаётся). Новых объектов схемы нет, временные таблицы уходят
-- по ON COMMIT DROP.
-- Прогон обеих фаз ДВАЖДЫ подряд с ROLLBACK вместо COMMIT — tradein-mvp/scripts/sql/
-- 286_dryrun.sql; второй проход обязан дать deleted=0 recomputed=0 nulled=0.
BEGIN;
-- Конвенция проекта (#2752): DELETE/UPDATE берут блокировки на строках
-- offer_price_history, а CREATE TEMP TABLE — DDL; без lock_timeout деплой встанет
-- в очередь за чужой сессией и утащит за собой запросы приложения.
SET LOCAL lock_timeout = '5s';
CREATE TEMP TABLE oph_decimal_slips ON COMMIT DROP AS
WITH pts AS (
SELECT id,
listing_id,
source,
price_rub,
lag(price_rub) OVER w AS prev_price,
lead(price_rub) OVER w AS next_price,
lead(id) OVER w AS next_id
FROM offer_price_history
WHERE change_time <> recorded_at
WINDOW w AS (PARTITION BY listing_id ORDER BY change_time, id)
CREATE TEMP TABLE oph_decimal_slips (
id bigint PRIMARY KEY,
listing_id bigint NOT NULL,
next_id bigint, -- строка, шедшая следом: у неё сменится база diff_percent
source text NOT NULL
) ON COMMIT DROP;
-- ФАЗА 1 — правило внутри серии, пошагово как в гейте.
DO $$
DECLARE
r record;
cur_listing bigint := NULL;
base numeric := NULL; -- последняя ОСТАВЛЕННАЯ цена этого объявления
witness numeric;
confirmed boolean;
BEGIN
FOR r IN
SELECT oph.id,
oph.listing_id,
oph.source,
oph.price_rub,
lead(oph.price_rub) OVER w AS next_price,
lead(oph.id) OVER w AS next_id,
l.price_rub AS listing_price
FROM offer_price_history oph
LEFT JOIN listings l ON l.id = oph.listing_id
WHERE oph.change_time <> oph.recorded_at
-- Объявления с одной строкой загрузчика кандидатов дать не могут ни по
-- одному из правил (нет ни базы слева, ни второй точки справа).
AND oph.listing_id IN (
SELECT listing_id
FROM offer_price_history
WHERE change_time <> recorded_at
GROUP BY listing_id
HAVING count(*) >= 2
)
WINDOW w AS (PARTITION BY oph.listing_id ORDER BY oph.change_time, oph.id)
ORDER BY oph.listing_id, oph.change_time, oph.id
LOOP
IF cur_listing IS DISTINCT FROM r.listing_id THEN
cur_listing := r.listing_id;
base := NULL; -- новое объявление, базы слева нет
END IF;
-- Свидетель справа: следующая точка серии, а у последней — текущая цена
-- объявления. Нет ни того, ни другого → witness IS NULL, и сравнение ниже
-- даёт NULL, то есть «не сдвиг»: без второго свидетеля точку не трогаем.
witness := COALESCE(r.next_price, r.listing_price);
IF base IS NOT NULL
AND base > 0 AND r.price_rub > 0 AND witness > 0
AND (r.price_rub / base BETWEEN 9.5 AND 10.5
OR base / r.price_rub BETWEEN 9.5 AND 10.5)
AND witness / base BETWEEN 0.9 AND 1.1
THEN
-- L-4: цену подтверждает триггерная строка того же объявления → это
-- наблюдавшаяся listings.price_rub, а не дефект разбора.
SELECT EXISTS (
SELECT 1
FROM offer_price_history t
WHERE t.listing_id = r.listing_id
AND t.change_time = t.recorded_at
AND t.price_rub = r.price_rub
)
INTO confirmed;
IF NOT confirmed THEN
INSERT INTO oph_decimal_slips (id, listing_id, next_id, source)
VALUES (r.id, r.listing_id, r.next_id, r.source);
CONTINUE; -- база НЕ двигается: этой точки в серии больше нет
END IF;
END IF;
IF r.price_rub > 0 THEN
base := r.price_rub;
END IF;
END LOOP;
END $$;
-- ФАЗА 2 — правило первой точки, по ОСТАВШИМСЯ строкам (кандидаты фазы 1 исключены).
-- Набор множественный, состояния не требует — обычный INSERT ... SELECT.
INSERT INTO oph_decimal_slips (id, listing_id, next_id, source)
WITH kept AS (
SELECT oph.id,
oph.listing_id,
oph.source,
oph.price_rub,
row_number() OVER w AS rn,
lead(oph.id) OVER w AS second_id,
lead(oph.price_rub) OVER w AS second_price,
lead(oph.price_rub, 2) OVER w AS third_price,
l.price_rub AS listing_price
FROM offer_price_history oph
LEFT JOIN listings l ON l.id = oph.listing_id
WHERE oph.change_time <> oph.recorded_at
AND NOT EXISTS (SELECT 1 FROM oph_decimal_slips s WHERE s.id = oph.id)
WINDOW w AS (PARTITION BY oph.listing_id ORDER BY oph.change_time, oph.id)
)
SELECT p.id,
p.listing_id,
p.source,
p.next_id,
-- Свидетель: следующая точка серии, а если её нет — текущая цена объявления.
COALESCE(p.next_price, l.price_rub) AS witness_price
FROM pts p
LEFT JOIN listings l ON l.id = p.listing_id
WHERE p.prev_price > 0
AND p.price_rub > 0
AND (p.price_rub / p.prev_price BETWEEN 9.5 AND 10.5
OR p.prev_price / p.price_rub BETWEEN 9.5 AND 10.5)
AND COALESCE(p.next_price, l.price_rub) > 0
AND COALESCE(p.next_price, l.price_rub) / p.prev_price BETWEEN 0.9 AND 1.1;
SELECT k.id, k.listing_id, k.second_id, k.source
FROM kept k
WHERE k.rn = 1
AND k.price_rub > 0
AND k.second_price > 0
AND (k.second_price / k.price_rub BETWEEN 9.5 AND 10.5
OR k.price_rub / k.second_price BETWEEN 9.5 AND 10.5)
AND COALESCE(k.third_price, k.listing_price) > 0
AND COALESCE(k.third_price, k.listing_price) / k.second_price BETWEEN 0.9 AND 1.1
AND NOT EXISTS (
SELECT 1
FROM offer_price_history t
WHERE t.listing_id = k.listing_id
AND t.change_time = t.recorded_at
AND t.price_rub = k.price_rub
);
DO $$
DECLARE
to_delete bigint;
per_source text;
deleted_rows bigint;
updated_rows bigint;
to_delete bigint;
per_source text;
deleted_rows bigint;
updated_rows bigint;
nulled_rows bigint;
nulled_src text;
BEGIN
SELECT count(*) INTO to_delete FROM oph_decimal_slips;
@ -110,7 +236,7 @@ BEGIN
IF to_delete > 200 THEN
RAISE EXCEPTION
'offer_price_history: к удалению % строк при замере 06.09.2026 = 76'
'offer_price_history: к удалению % строк при предзамере 06.09.2026 ≈ 47'
'критерий ловит не то, миграция остановлена', to_delete;
END IF;
@ -150,6 +276,27 @@ BEGIN
RAISE NOTICE 'offer_price_history: удалено строк = %, пересчитано diff_percent у соседей = %',
deleted_rows, updated_rows;
-- Финальный шаг: величина не того рода → NULL. Правило validate_diff_percent,
-- применённое ко ВСЕМ источникам и только к строкам загрузчика.
SELECT string_agg(source || '=' || cnt, ', ' ORDER BY source)
INTO nulled_src
FROM (
SELECT source, count(*) AS cnt
FROM offer_price_history
WHERE change_time <> recorded_at
AND abs(diff_percent) > 100
GROUP BY source
) s;
UPDATE offer_price_history
SET diff_percent = NULL
WHERE change_time <> recorded_at
AND abs(diff_percent) > 100;
GET DIAGNOSTICS nulled_rows = ROW_COUNT;
RAISE NOTICE 'offer_price_history: |diff_percent| > 100 обнулено = % (%)',
nulled_rows, COALESCE(nulled_src, 'ни одной');
END $$;
COMMIT;

View file

@ -32,6 +32,7 @@ importers) — `scraper_kit.domclick_exceptions` остаётся единств
from __future__ import annotations
import itertools
import json
import logging
from datetime import UTC, datetime, timedelta
@ -41,7 +42,11 @@ import httpx
import pytest
from scraper_kit.browser_fetcher import SidecarBanPageError
from scraper_kit.domclick_exceptions import DomClickBlockedError, DomClickParseError
from scraper_kit.offer_price_history import drop_decimal_slips, validate_diff_percent
from scraper_kit.offer_price_history import (
drop_decimal_slips,
recompute_diff_percent,
validate_diff_percent,
)
from scraper_kit.providers.domclick.detail import (
DomClickDetailEnrichment,
_extract_ssr_state,
@ -761,8 +766,15 @@ def test_drop_decimal_slips_keeps_honest_doubling() -> None:
def test_drop_decimal_slips_keeps_tenfold_without_return_to_base() -> None:
"""×10 БЕЗ возврата к базе — возможное честное изменение цены, точку не трогаем."""
series = [1020000, 10200000, 10150000]
"""×10 БЕЗ возврата к базе — возможное честное изменение цены, точку не трогаем.
Хвост 1 020 000 10 200 000 10 150 000 намеренно стоит НЕ первым: ровно та
же тройка в начале серии это уже сдвиг по правилу первой точки (см.
test_drop_decimal_slips_drops_first_point_witnessed_by_third), и правила здесь
не спорят, а смотрят на разные концы серии. Внутри серии у 10 200 000 есть
база слева, и она НЕ подтверждена следующей точкой (9.95 1) точка честная.
"""
series = [1000000, 1020000, 10200000, 10150000]
assert drop_decimal_slips(series, lambda p: p) == (series, 0)
@ -779,6 +791,136 @@ def test_drop_decimal_slips_last_point_kept_without_current_price() -> None:
assert drop_decimal_slips(series, lambda p: p, current_price=None) == (series, 0)
# ── правило ПЕРВОЙ точки (#3376, прод-разбор 06.09.2026) ─────────────────────
# 12 из 20 остатков domklik — серии, у которых дефектная точка первая: 330 000 →
# 3 300 000 (текущая цена 3 300 000), 420 000 → 4 200 000 → 4 500 000. Слева базы
# нет, основное правило такую точку не видит. Свидетелей по-прежнему двое.
def test_drop_decimal_slips_drops_first_point_witnessed_by_current_price() -> None:
"""330 000 → 3 300 000 при цене 3 300 000: вторую точку подтверждает цена лота."""
kept, dropped = drop_decimal_slips([330000, 3300000], lambda p: p, current_price=3300000)
assert (kept, dropped) == ([3300000], 1)
def test_drop_decimal_slips_drops_first_point_witnessed_by_third() -> None:
"""420 000 → 4 200 000 → 4 500 000: вторую точку подтверждает третья."""
kept, dropped = drop_decimal_slips([420000, 4200000, 4500000], lambda p: p)
assert (kept, dropped) == ([4200000, 4500000], 1)
def test_drop_decimal_slips_keeps_first_point_when_series_runs_away() -> None:
"""Первая ×10 ко второй, но третья ушла ЕЩЁ дальше — это разгон цены, не разряд."""
series = [400000, 4000000, 8000000]
assert drop_decimal_slips(series, lambda p: p) == (series, 0)
def test_drop_decimal_slips_keeps_first_point_without_second_witness() -> None:
"""Две точки и текущей цены нет — второго свидетеля не существует."""
series = [330000, 3300000]
assert drop_decimal_slips(series, lambda p: p, current_price=None) == (series, 0)
def test_drop_decimal_slips_first_point_decided_on_kept_series() -> None:
"""Свидетели первой точки — ОСТАВШИЕСЯ соседи, а не сырые.
1M 10M 100M 10M: 100M выбрасывает основное правило, и только после этого
видно, что первая точка ÷10 к оставшейся серии [10M, 10M]. Считай правило по
сырым соседям первый проход оставил бы 1M, а второй выбросил, то есть гейт
перестал бы быть идемпотентным (перебор ниже ловит 1122 такие серии).
"""
kept, dropped = drop_decimal_slips(
[1_000_000, 10_000_000, 100_000_000, 10_000_000], lambda p: p
)
assert (kept, dropped) == ([10_000_000, 10_000_000], 2)
# ── база = предыдущая ОСТАВЛЕННАЯ точка (фальсификаторы ревьюера) ────────────
def test_drop_decimal_slips_alternating_series_keeps_honest_points() -> None:
"""1M → 10M → 1M → 10M при цене 1M: обе честные точки на месте.
Первая редакция миграции 286 брала базой СЫРУЮ предыдущую строку (lag) и на
этой серии оставляла [1M] вместо [1M, 1M] удаляла честную точку.
"""
kept, dropped = drop_decimal_slips(
[1_000_000, 10_000_000, 1_000_000, 10_000_000], lambda p: p, current_price=1_000_000
)
assert (kept, dropped) == ([1_000_000, 1_000_000], 2)
def test_drop_decimal_slips_base_is_previous_kept_not_raw() -> None:
"""1M → 10M → 1.05M → 9.9M при цене 1M: после выброса 10M база для 1.05M — это 1M.
По сырой базе 1.05M тоже стала бы кандидатом (10M/1.05M = 9.52), а следующий
прогон доел бы 9.9M ровно та неидемпотентность, которую нашёл ревьюер.
"""
kept, dropped = drop_decimal_slips(
[1_000_000, 10_000_000, 1_050_000, 9_900_000], lambda p: p, current_price=1_000_000
)
assert (kept, dropped) == ([1_000_000, 1_050_000, 9_900_000], 1)
# ── свойство: gate(gate(s)) == gate(s) ───────────────────────────────────────
_SLIP_ALPHABET = (1_000_000, 10_000_000, 100_000_000, 1_050_000, 9_900_000)
_SLIP_CURRENTS = (None, 1_000_000, 10_000_000)
def test_drop_decimal_slips_is_idempotent(caplog: pytest.LogCaptureFixture) -> None:
"""Повторный прогон гейта не трогает ничего — на всех сериях длины 2-6.
Свойство несущее, а не декоративное: миграция 286 зеркалит эту функцию
пошагово, а доказать идемпотентность на самом SQL-файле нечем прод-прогон
один. Раз выборка кандидатов повторяет гейт, свойство переносится на неё.
ФАЛЬСИФИЦИРУЕМОСТЬ (замерено на этом же переборе 19 525 серий × 3 цены):
база = сырая предыдущая точка вместо оставленной 136 неидемпотентных
прогонов, первый же (1M, 10M, 1.05M, 9.9M) при цене 1M: [1M, 9.9M] [1M];
решение по первой точке на сырых соседях вместо оставшихся 1122 прогона,
первый (1M, 10M, 100M) при цене 10M: [1M, 10M] [10M].
"""
# Гейт логирует каждый выброс; на таком переборе это десятки тысяч записей.
caplog.set_level(logging.CRITICAL, logger="scraper_kit.offer_price_history")
for length in range(2, 7):
for series in itertools.product(_SLIP_ALPHABET, repeat=length):
for current in _SLIP_CURRENTS:
once, _ = drop_decimal_slips(list(series), lambda p: p, current_price=current)
twice, dropped_again = drop_decimal_slips(once, lambda p: p, current_price=current)
assert (twice, dropped_again) == (once, 0), (
f"серия {series} при текущей цене {current}: {once}{twice}"
)
# ── recompute_diff_percent — общий пересчёт после выброса ────────────────────
def test_recompute_diff_percent_from_surviving_neighbours() -> None:
"""Процент считается от оставшейся базы; у первой точки базы нет → NULL."""
changes: list[dict[str, object]] = [
{"price_rub": 3720000, "diff_percent": -89.87},
{"price_rub": 3770000, "diff_percent": 900.0},
]
recompute_diff_percent(changes)
assert [c["diff_percent"] for c in changes] == [None, 1.34]
def test_recompute_diff_percent_survives_none_price() -> None:
"""price_rub=None не роняет пересчёт: ручной ingest кладёт price_changes сырьём.
scripts/ingest_domclick_jsonl.py:108 берёт rec["price_changes"] из JSONL без
валидации арифметика по None здесь дала бы TypeError на всю запись.
"""
changes: list[dict[str, object]] = [
{"price_rub": 5_000_000, "diff_percent": 12.0},
{"price_rub": None, "diff_percent": 900.0},
{"price_rub": 5_100_000, "diff_percent": -90.0},
]
recompute_diff_percent(changes)
assert [c["diff_percent"] for c in changes] == [None, None, 2.0]
def test_save_detail_enrichment_drops_decimal_slip_before_insert() -> None:
"""Писатель не отправляет спайк в INSERT, а соседу пересчитывает diff_percent."""
db = MagicMock()
@ -806,6 +948,36 @@ def test_save_detail_enrichment_drops_decimal_slip_before_insert() -> None:
assert insert_calls[1][0][1]["diff"] == 1.34
def test_save_detail_enrichment_witness_comes_only_from_listing_price() -> None:
"""Серия из ДВУХ точек: второго свидетеля даёт исключительно цена листинга.
Здесь чтение цены не декорация, а единственный вход гейта: у последней точки
следующей нет. Сломай его (RETURNING убран / fetchone None) свидетеля не
станет, спайк 10 200 000 уедет в offer_price_history, и тест покраснеет на
втором INSERT'е. Проверено подменой fetchone → None.
"""
db = MagicMock()
db.execute.return_value.rowcount = 1
db.execute.return_value.fetchone.return_value = MagicMock(price_rub=1020000)
times = [datetime(2026, 4, d, tzinfo=UTC) for d in (1, 2)]
e = DomClickDetailEnrichment(
item_id="x",
source_url=_CARD_URL,
price_changes=[
{"change_time": times[0], "price_rub": 1020000, "diff_percent": None},
{"change_time": times[1], "price_rub": 10200000, "diff_percent": 900.0},
],
)
assert save_detail_enrichment(db, 999, e) is True
insert_calls = [c for c in db.execute.call_args_list if "offer_price_history" in str(c[0][0])]
assert [c[0][1]["price"] for c in insert_calls] == [1020000]
# Цену отдаёт тот же UPDATE — отдельного SELECT по PK больше нет.
assert "RETURNING price_rub" in str(db.execute.call_args_list[0][0][0])
assert not [c for c in db.execute.call_args_list if "SELECT price_rub" in str(c[0][0])]
# ── canon_sale_type (#2674) ───────────────────────────────────────────────────
# Кейсы — ФАКТИЧЕСКИЙ словарь прода на 2026-08-06:
# SELECT source, sale_type, count(*) FROM listings GROUP BY 1,2 →

View file

@ -62,9 +62,10 @@ class _FakeNested:
class _FakeSession:
"""Пишет все execute(stmt, params); rowcount=1 у каждого оператора.
fetchone() None: writer читает текущую цену листинга для гейта сдвигов разряда
(#3376), а этой фикстуре цена не нужна — «строки нет» гейт трактует как отсутствие
второго свидетеля и серию не трогает.
fetchone() None: UPDATE listings отдаёт текущую цену через RETURNING она нужна
гейту сдвигов разряда (#3376) как второй свидетель. Этой фикстуре цена не нужна,
«строки нет» гейт трактует как отсутствие свидетеля и серию не трогает; факт
«листинг найден» здесь по-прежнему приходит из rowcount, а не из RETURNING.
"""
def __init__(self) -> None:

View file

@ -308,6 +308,34 @@ def test_save_detail_enrichment_oph_on_conflict_constraint():
)
def test_save_detail_enrichment_recomputes_diff_after_decimal_slip():
"""cian: после выброса сдвига разряда соседу пересчитывается diff_percent (#3376).
377 000 3 770 000 377 000 при цене листинга 377 000: средняя точка ×10 с
возвратом к базе выбрасывается. У последней точки парсер посчитал 90% от
исчезнувшей базы 3 770 000; от настоящей базы 377 000 это 0%. Без пересчёта в
колонку уехал бы процент от цены, которой в истории больше нет (у domclick это
уже чинилось, cian отставал).
"""
db = _mock_db_detail(price_rub=377_000)
enrichment = DetailEnrichment(
price_changes=[
{"change_time": "2026-04-01T00:00:00Z", "price_rub": 377_000, "diff_percent": None},
{"change_time": "2026-04-05T00:00:00Z", "price_rub": 3_770_000, "diff_percent": 900.0},
{"change_time": "2026-04-10T00:00:00Z", "price_rub": 377_000, "diff_percent": -90.0},
]
)
save_detail_enrichment(db, 23, enrichment)
oph = [
(params["price"], params["diff"])
for sql, params, _ in db._execute_log
if "offer_price_history" in sql
]
assert oph == [(377_000, None), (377_000, 0.0)]
def test_save_detail_enrichment_skips_price_change_without_change_time():
"""price_changes без change_time пропускаются, не вызывают INSERT."""
db = _mock_db_detail(price_rub=5_000_000)

View file

@ -14,12 +14,18 @@ loader клал поле источника ``diff``, а там рубли. Пр
Цена решения: настоящий рост цены больше чем в 2 раза (+100%) тоже уйдёт в NULL.
Это осознанный размен колонка аналитическая, а молчаливый мусор в ней дороже
пропущенного выброса.
Здесь же живёт второй гейт той же границы записи drop_decimal_slips (#3376):
он выбрасывает из СЕРИИ точки с потерянным разрядом источника (×10/÷10 с
подтверждением у соседей), а recompute_diff_percent чинит процент у соседа
выброшенной точки. Оба вызываются писателями истории (domclick/detail.py,
cian/detail.py) и зеркалятся миграцией 286.
"""
from __future__ import annotations
import logging
from collections.abc import Callable, Sequence
from collections.abc import Callable, MutableMapping, Sequence
from typing import TypeVar
logger = logging.getLogger(__name__)
@ -78,19 +84,37 @@ def _positive_price(value: object) -> float | None:
return numeric if numeric > 0 else None
def _shifted_by_decade(a: float, b: float) -> bool:
"""b отличается от a ровно на разряд — ×10 или ÷10 (окно 9.5..10.5)."""
return (_SLIP_RATIO_MIN <= b / a <= _SLIP_RATIO_MAX) or (
_SLIP_RATIO_MIN <= a / b <= _SLIP_RATIO_MAX
)
def _is_decimal_slip(prev: float, cur: float, nxt: float | None) -> bool:
"""Точка = сдвиг разряда: ×10 (или ÷10) к prev И возврат к базе у next."""
ratio = cur / prev
shifted = (_SLIP_RATIO_MIN <= ratio <= _SLIP_RATIO_MAX) or (
1 / _SLIP_RATIO_MAX <= ratio <= 1 / _SLIP_RATIO_MIN
)
if not shifted or nxt is None:
if not _shifted_by_decade(prev, cur) or nxt is None:
# nxt is None — второго свидетеля нет, и это НЕ повод выбрасывать: ×10 без
# возврата к базе бывает честным изменением цены.
return False
return _BASE_RETURN_MIN <= nxt / prev <= _BASE_RETURN_MAX
def _is_first_point_slip(first: float, second: float, witness: float | None) -> bool:
"""Первая точка серии = сдвиг разряда: у неё нет prev, свидетели справа.
Прод-разбор 06.09.2026: 12 из 20 остатков domklik серии вида
``330 000 3 300 000`` (текущая цена 3 300 000) и ``420 000 4 200 000
4 500 000``. Дефектная точка ПЕРВАЯ, основное правило её не видит: базы слева
нет. Свидетелей по-прежнему двое ×10 ко ВТОРОЙ точке и подтверждение второй
точки третьей-или-текущей-ценой. Одного скачка мало: ``400 000 4 000 000
8 000 000`` это разгон цены, а не потерянный разряд.
"""
if not _shifted_by_decade(first, second) or witness is None:
return False
return _BASE_RETURN_MIN <= witness / second <= _BASE_RETURN_MAX
def drop_decimal_slips(
points: Sequence[T],
price_of: Callable[[T], object],
@ -108,6 +132,22 @@ def drop_decimal_slips(
Критерий требует ДВУХ свидетелей: скачок ×10 к предыдущей точке И возврат к
базе у следующей. Одного скачка мало так выглядит и честная смена цены.
У ПЕРВОЙ точки предыдущей нет, для неё работает зеркальное правило
(_is_first_point_slip) оба свидетеля справа.
ИДЕМПОТЕНТНОСТЬ свойство, а не пожелание: ``gate(gate(s)) == gate(s)``.
Его держит property-тест перебором серий длины 2-6 (см.
test_drop_decimal_slips_is_idempotent), и от него же зависит миграция 286,
которая эту функцию зеркалит: миграция обязана быть безопасной при повторном
прогоне. Два места, где свойство легко потерять:
база предыдущая ОСТАВЛЕННАЯ точка, а не предыдущая сырая. Иначе на серии
1M 10M 1.05M 9.9M первый проход выбросит и 1.05M (её база уже
выброшенная 10M), а второй доест 9.9M.
решение по ПЕРВОЙ точке принимается по ОСТАВШЕЙСЯ серии, а не по сырой:
свидетели первой точки (вторая и третья точки) сами могут оказаться
сдвигами. На сырой серии 1M 10M 100M 10M первый проход оставил бы
первую точку (третья сырая точка 100M второй не подтверждает), а второй
выбросил, потому что 100M к тому моменту уже удалена.
Args:
points: точки одного объявления, упорядоченные по времени (возрастание).
@ -120,7 +160,9 @@ def drop_decimal_slips(
(серия без сдвигов, число выброшенных точек).
"""
prices = [_positive_price(price_of(point)) for point in points]
now_price = _positive_price(current_price)
kept: list[T] = []
kept_prices: list[float | None] = []
dropped = 0
prev: float | None = None
for i, point in enumerate(points):
@ -128,7 +170,7 @@ def drop_decimal_slips(
# База — предыдущая ОСТАВЛЕННАЯ точка: у двух сдвигов подряд второй иначе
# считался бы от выброшенного соседа.
if prev is not None and cur is not None:
nxt = prices[i + 1] if i + 1 < len(prices) else _positive_price(current_price)
nxt = prices[i + 1] if i + 1 < len(prices) else now_price
if _is_decimal_slip(prev, cur, nxt):
dropped += 1
logger.warning(
@ -141,6 +183,42 @@ def drop_decimal_slips(
)
continue
kept.append(point)
kept_prices.append(cur)
if cur is not None:
prev = cur
# Первая точка — последней: её свидетели это уже ОСТАВШИЕСЯ вторая и третья
# точки (см. про идемпотентность выше).
if len(kept_prices) >= 2 and kept_prices[0] is not None and kept_prices[1] is not None:
witness = kept_prices[2] if len(kept_prices) >= 3 else now_price
if _is_first_point_slip(kept_prices[0], kept_prices[1], witness):
logger.warning(
"offer_price_history: первая точка %s выброшена как сдвиг разряда "
"(listing_id=%s, next=%s, witness=%s)",
kept_prices[0],
listing_id,
kept_prices[1],
witness,
)
return kept[1:], dropped + 1
return kept, dropped
def recompute_diff_percent(changes: Sequence[MutableMapping[str, object]]) -> None:
"""Пересчитать diff_percent из соседних цен ПО МЕСТУ — после выброса точек.
Формула та же, что у починенного парсера (#3225) и у миграции 285:
(price prev)/prev*100, prev нет NULL. Ноль здесь читался бы как «цена не
менялась», то есть как измерение, которого не было.
Точки с непригодной ценой (None из ручного ingest, 0, мусор) базой не
становятся и получают NULL: писатели их всё равно не вставляют.
"""
prev: float | None = None
for change in changes:
price = _positive_price(change.get("price_rub"))
if price is None:
change["diff_percent"] = None
continue
change["diff_percent"] = round((price - prev) / prev * 100, 2) if prev else None
prev = price

View file

@ -25,7 +25,11 @@ from sqlalchemy.orm import Session
from scraper_kit.ceiling_height import plausible_ceiling_m
from scraper_kit.cian_exceptions import CianBlockedError
from scraper_kit.cian_state_parser import extract_all_states, extract_state
from scraper_kit.offer_price_history import drop_decimal_slips, validate_diff_percent
from scraper_kit.offer_price_history import (
drop_decimal_slips,
recompute_diff_percent,
validate_diff_percent,
)
from scraper_kit.providers._base import build_curl_cffi_session
from scraper_kit.providers._proxy import curl_proxy_url
from scraper_kit.repair_state_normalizer import (
@ -481,6 +485,11 @@ def save_detail_enrichment(
current_price=_snap_row.price_rub if _snap_row is not None else None,
listing_id=listing_id,
)
if decimal_slips_dropped:
# То же, что у domclick: diff_percent соседа считался от базы, которой
# больше нет. Без пересчёта у него остаётся процент от фантомной цены
# (377 000 → 3 770 000 → 377 000 оставлял бы у последней точки 90%).
recompute_diff_percent(_changes)
for change in _changes:
db.execute(
text("""

View file

@ -62,7 +62,11 @@ from scraper_kit.domclick_exceptions import (
DomClickParseError,
)
from scraper_kit.browser_fetcher import SidecarBanPageError
from scraper_kit.offer_price_history import drop_decimal_slips, validate_diff_percent
from scraper_kit.offer_price_history import (
drop_decimal_slips,
recompute_diff_percent,
validate_diff_percent,
)
from scraper_kit.repair_state_normalizer import (
infer_repair_state_from_text,
normalize_repair_state,
@ -774,6 +778,7 @@ def save_detail_enrichment(
raw_payload = COALESCE(raw_payload, '{}'::jsonb) || CAST(:raw_extra AS jsonb),
detail_enriched_at = NOW()
WHERE id = CAST(:lid AS bigint)
RETURNING price_rub
"""),
{
"lid": listing_id,
@ -791,32 +796,28 @@ def save_detail_enrichment(
"raw_extra": raw_extra_json,
},
)
# rowcount читаем ДО fetchone(): у RETURNING-результата курсор после выборки
# уже мягко закрыт, и порядок «сначала счётчик, потом строка» снимает вопрос,
# доживёт ли rowcount до этого момента у конкретного драйвера.
found = result.rowcount > 0
# Текущая цена листинга — второй свидетель для последней точки priceHistory
# (#3376). Берём RETURNING'ом у UPDATE выше, а не отдельным SELECT по PK:
# это та же строка, только что обновлённая, лишний round-trip не нужен.
listing_row = result.fetchone()
# Сдвиги разряда (#3376) — выбрасываем ДО вставки. Серия уже отсортирована по
# change_time парсером. Текущая цена листинга — второй свидетель для последней
# точки; detail-backfill цену не меняет, так что строка listings актуальна.
current_price_row = db.execute(
text("SELECT price_rub FROM listings WHERE id = CAST(:lid AS bigint)"),
{"lid": listing_id},
).fetchone()
# change_time парсером.
price_changes, decimal_slips_dropped = drop_decimal_slips(
e.price_changes,
lambda c: c.get("price_rub"),
current_price=current_price_row.price_rub if current_price_row else None,
current_price=listing_row.price_rub if listing_row else None,
listing_id=listing_id,
)
if decimal_slips_dropped:
# diff_percent парсер считал от соседа, которого больше нет (#3225: процент
# берётся из соседних цен, а не из поля источника). Пересчитываем по той же
# формуле, иначе у следующей точки останется процент от фантомной базы.
prev_price: int | None = None
for change in price_changes:
change["diff_percent"] = (
round((change["price_rub"] - prev_price) / prev_price * 100, 2)
if prev_price
else None
)
prev_price = change["price_rub"]
recompute_diff_percent(price_changes)
# priceHistory — INSERT, де-дуп по (listing_id, change_time). SAVEPOINT per-row
# (backend.md): битый cast одной записи не валит весь UPDATE+commit.
@ -857,11 +858,10 @@ def save_detail_enrichment(
)
# Дом-поля — только если листинг нашёлся: без строки в listings нет и house_id_fk.
if result.rowcount > 0:
if found:
_fill_house_params_from_detail(db, listing_id, e)
db.commit()
found = result.rowcount > 0
if not found:
logger.warning(
"domclick_detail save: listing not found id=%s (item_id=%s)",

View file

@ -0,0 +1,252 @@
-- 286_dryrun.sql — прогон миграции 286 на проде БЕЗ записи (#3376).
--
-- ЗАЧЕМ. Миграция 286 удаляет строки из offer_price_history. Прежде чем этому
-- уехать на деплой, надо увидеть на БОЕВЫХ данных: сколько строк она заберёт по
-- источникам, сколько соседей пересчитает, сколько процентов обнулит — и что
-- ВТОРОЙ прогон подряд не делает уже ничего (идемпотентность).
--
-- КАК ЗАПУСКАТЬ (ничего не остаётся, в конце ROLLBACK):
-- ssh poincare
-- docker exec -i tradein-postgres psql -U <user> -d tradein -v ON_ERROR_STOP=1 \
-- -f - < 286_dryrun.sql
-- (или скопировать файл в контейнер и psql -f). Смотреть строки NOTICE.
--
-- ЧЕМ ОТЛИЧАЕТСЯ ОТ 286 — ровно тремя вещами, всё остальное скопировано дословно:
-- 1. вся работа обёрнута в FOR pass_no IN 1..2 — два прохода подряд в ОДНОЙ
-- транзакции; второй обязан дать deleted=0 recomputed=0 nulled=0;
-- 2. порог 200 даёт WARNING, а не EXCEPTION: на dry-run важнее увидеть число и
-- добежать до второго прохода, чем оборвать транзакцию;
-- 3. ROLLBACK вместо COMMIT.
-- ФАЙЛ РАЗОВЫЙ и живёт вне data/sql (автоприменение его не подхватывает). Правишь
-- 286 — правь и здесь, иначе сверять будет нечего.
BEGIN;
SET LOCAL lock_timeout = '5s';
-- ── ПРЕДЗАМЕР ОДНИМ SELECT (для сверки с числами в шапке 286) ────────────────
-- Правило внутри серии пошаговое (база = предыдущая ОСТАВЛЕННАЯ точка) и одним
-- SELECT не выражается — его считает DO-блок ниже. А эти два выражаются.
-- (а) кандидаты правила ПЕРВОЙ точки по СЫРЫМ строкам, до удалений фазы 1.
-- Ожидание из шапки 286: ~12, из них домклик большинство.
SELECT 'предзамер: первая точка (сырые строки)' AS metric, k.source, count(*)
FROM (
SELECT oph.source,
oph.price_rub,
row_number() OVER w AS rn,
lead(oph.price_rub) OVER w AS second_price,
lead(oph.price_rub, 2) OVER w AS third_price,
l.price_rub AS listing_price
FROM offer_price_history oph
LEFT JOIN listings l ON l.id = oph.listing_id
WHERE oph.change_time <> oph.recorded_at
WINDOW w AS (PARTITION BY oph.listing_id ORDER BY oph.change_time, oph.id)
) k
WHERE k.rn = 1
AND k.price_rub > 0
AND k.second_price > 0
AND (k.second_price / k.price_rub BETWEEN 9.5 AND 10.5
OR k.price_rub / k.second_price BETWEEN 9.5 AND 10.5)
AND COALESCE(k.third_price, k.listing_price) > 0
AND COALESCE(k.third_price, k.listing_price) / k.second_price BETWEEN 0.9 AND 1.1
GROUP BY k.source
ORDER BY k.source;
-- (б) сколько строк заберёт финальный шаг |diff_percent| > 100 → NULL.
SELECT 'предзамер: |diff_percent| > 100' AS metric, source, count(*)
FROM offer_price_history
WHERE change_time <> recorded_at
AND abs(diff_percent) > 100
GROUP BY source
ORDER BY source;
-- ── ДВА ПРОХОДА ЛОГИКИ 286 ───────────────────────────────────────────────────
CREATE TEMP TABLE oph_decimal_slips (
id bigint PRIMARY KEY,
listing_id bigint NOT NULL,
next_id bigint,
source text NOT NULL
) ON COMMIT DROP;
DO $$
DECLARE
pass_no int;
r record;
cur_listing bigint;
base numeric;
witness numeric;
confirmed boolean;
to_delete bigint;
per_source text;
deleted_rows bigint;
updated_rows bigint;
nulled_rows bigint;
nulled_src text;
BEGIN
FOR pass_no IN 1..2 LOOP
DELETE FROM oph_decimal_slips;
cur_listing := NULL;
base := NULL;
-- ФАЗА 1 — правило внутри серии, пошагово как в гейте.
FOR r IN
SELECT oph.id,
oph.listing_id,
oph.source,
oph.price_rub,
lead(oph.price_rub) OVER w AS next_price,
lead(oph.id) OVER w AS next_id,
l.price_rub AS listing_price
FROM offer_price_history oph
LEFT JOIN listings l ON l.id = oph.listing_id
WHERE oph.change_time <> oph.recorded_at
AND oph.listing_id IN (
SELECT listing_id
FROM offer_price_history
WHERE change_time <> recorded_at
GROUP BY listing_id
HAVING count(*) >= 2
)
WINDOW w AS (PARTITION BY oph.listing_id ORDER BY oph.change_time, oph.id)
ORDER BY oph.listing_id, oph.change_time, oph.id
LOOP
IF cur_listing IS DISTINCT FROM r.listing_id THEN
cur_listing := r.listing_id;
base := NULL;
END IF;
witness := COALESCE(r.next_price, r.listing_price);
IF base IS NOT NULL
AND base > 0 AND r.price_rub > 0 AND witness > 0
AND (r.price_rub / base BETWEEN 9.5 AND 10.5
OR base / r.price_rub BETWEEN 9.5 AND 10.5)
AND witness / base BETWEEN 0.9 AND 1.1
THEN
SELECT EXISTS (
SELECT 1
FROM offer_price_history t
WHERE t.listing_id = r.listing_id
AND t.change_time = t.recorded_at
AND t.price_rub = r.price_rub
)
INTO confirmed;
IF NOT confirmed THEN
INSERT INTO oph_decimal_slips (id, listing_id, next_id, source)
VALUES (r.id, r.listing_id, r.next_id, r.source);
CONTINUE;
END IF;
END IF;
IF r.price_rub > 0 THEN
base := r.price_rub;
END IF;
END LOOP;
-- ФАЗА 2 — правило первой точки, по оставшимся строкам.
INSERT INTO oph_decimal_slips (id, listing_id, next_id, source)
WITH kept AS (
SELECT oph.id,
oph.listing_id,
oph.source,
oph.price_rub,
row_number() OVER w AS rn,
lead(oph.id) OVER w AS second_id,
lead(oph.price_rub) OVER w AS second_price,
lead(oph.price_rub, 2) OVER w AS third_price,
l.price_rub AS listing_price
FROM offer_price_history oph
LEFT JOIN listings l ON l.id = oph.listing_id
WHERE oph.change_time <> oph.recorded_at
AND NOT EXISTS (SELECT 1 FROM oph_decimal_slips s WHERE s.id = oph.id)
WINDOW w AS (PARTITION BY oph.listing_id ORDER BY oph.change_time, oph.id)
)
SELECT k.id, k.listing_id, k.second_id, k.source
FROM kept k
WHERE k.rn = 1
AND k.price_rub > 0
AND k.second_price > 0
AND (k.second_price / k.price_rub BETWEEN 9.5 AND 10.5
OR k.price_rub / k.second_price BETWEEN 9.5 AND 10.5)
AND COALESCE(k.third_price, k.listing_price) > 0
AND COALESCE(k.third_price, k.listing_price) / k.second_price BETWEEN 0.9 AND 1.1
AND NOT EXISTS (
SELECT 1
FROM offer_price_history t
WHERE t.listing_id = k.listing_id
AND t.change_time = t.recorded_at
AND t.price_rub = k.price_rub
);
SELECT count(*) INTO to_delete FROM oph_decimal_slips;
SELECT string_agg(source || '=' || cnt, ', ' ORDER BY source)
INTO per_source
FROM (
SELECT source, count(*) AS cnt
FROM oph_decimal_slips
GROUP BY source
) s;
-- Отличие 2 от 286: здесь WARNING, а не EXCEPTION.
IF to_delete > 200 THEN
RAISE WARNING 'pass%: кандидатов % > порога 200 — на проде 286 здесь УПАЛА БЫ',
pass_no, to_delete;
END IF;
DELETE FROM offer_price_history oph
USING oph_decimal_slips s
WHERE oph.id = s.id;
GET DIAGNOSTICS deleted_rows = ROW_COUNT;
WITH neighbours AS (
SELECT id,
price_rub,
lag(price_rub) OVER (PARTITION BY listing_id ORDER BY change_time, id)
AS prev_price
FROM offer_price_history
WHERE change_time <> recorded_at
AND listing_id IN (SELECT DISTINCT listing_id FROM oph_decimal_slips)
),
recomputed AS (
SELECT n.id,
CASE
WHEN n.prev_price > 0
AND abs((n.price_rub - n.prev_price) / n.prev_price * 100) <= 100
THEN round((n.price_rub - n.prev_price) / n.prev_price * 100, 2)
END AS new_diff
FROM neighbours n
WHERE n.id IN (SELECT next_id FROM oph_decimal_slips WHERE next_id IS NOT NULL)
)
UPDATE offer_price_history oph
SET diff_percent = rc.new_diff
FROM recomputed rc
WHERE oph.id = rc.id
AND oph.diff_percent IS DISTINCT FROM rc.new_diff;
GET DIAGNOSTICS updated_rows = ROW_COUNT;
SELECT string_agg(source || '=' || cnt, ', ' ORDER BY source)
INTO nulled_src
FROM (
SELECT source, count(*) AS cnt
FROM offer_price_history
WHERE change_time <> recorded_at
AND abs(diff_percent) > 100
GROUP BY source
) s;
UPDATE offer_price_history
SET diff_percent = NULL
WHERE change_time <> recorded_at
AND abs(diff_percent) > 100;
GET DIAGNOSTICS nulled_rows = ROW_COUNT;
RAISE NOTICE 'pass%: deleted=% (%) recomputed=% nulled=% (%)',
pass_no,
deleted_rows, COALESCE(per_source, 'ни одного'),
updated_rows,
nulled_rows, COALESCE(nulled_src, 'ни одной');
END LOOP;
END $$;
ROLLBACK;