From 721ceb987655b9ad0bd14ceabc0df3a6d33e8ca3 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sun, 6 Sep 2026 03:53:08 +0500 Subject: [PATCH] =?UTF-8?q?fix(tradein):=20=D0=B2=D1=8B=D0=B1=D0=BE=D1=80?= =?UTF-8?q?=D0=BA=D0=B0=20=D0=BC=D0=B8=D0=B3=D1=80=D0=B0=D1=86=D0=B8=D0=B8?= =?UTF-8?q?=20286=20=D0=BF=D0=BE=D0=B2=D1=82=D0=BE=D1=80=D1=8F=D0=B5=D1=82?= =?UTF-8?q?=20=D0=B3=D0=B5=D0=B9=D1=82=201:1,=20=D0=BF=D1=80=D0=B0=D0=B2?= =?UTF-8?q?=D0=B8=D0=BB=D0=BE=20=D0=BF=D0=B5=D1=80=D0=B2=D0=BE=D0=B9=20?= =?UTF-8?q?=D1=82=D0=BE=D1=87=D0=BA=D0=B8=20(#3376)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Ревью нашло, что миграция и код ловили РАЗНОЕ. Миграция брала базой предыдущую СЫРУЮ строку (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 без валидации). --- .../app/services/yandex_price_history.py | 68 +++-- .../286_offer_price_history_decimal_slips.sql | 265 ++++++++++++++---- .../tests/scrapers/test_domclick_detail.py | 178 +++++++++++- .../tests/test_3253_domclick_house_fields.py | 7 +- .../backend/tests/test_snapshot_writer.py | 28 ++ .../src/scraper_kit/offer_price_history.py | 92 +++++- .../src/scraper_kit/providers/cian/detail.py | 11 +- .../scraper_kit/providers/domclick/detail.py | 36 +-- tradein-mvp/scripts/sql/286_dryrun.sql | 252 +++++++++++++++++ 9 files changed, 810 insertions(+), 127 deletions(-) create mode 100644 tradein-mvp/scripts/sql/286_dryrun.sql diff --git a/tradein-mvp/backend/app/services/yandex_price_history.py b/tradein-mvp/backend/app/services/yandex_price_history.py index 57d74ecb..17ff09a4 100644 --- a/tradein-mvp/backend/app/services/yandex_price_history.py +++ b/tradein-mvp/backend/app/services/yandex_price_history.py @@ -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 diff --git a/tradein-mvp/backend/data/sql/286_offer_price_history_decimal_slips.sql b/tradein-mvp/backend/data/sql/286_offer_price_history_decimal_slips.sql index 7b07557e..d6332393 100644 --- a/tradein-mvp/backend/data/sql/286_offer_price_history_decimal_slips.sql +++ b/tradein-mvp/backend/data/sql/286_offer_price_history_decimal_slips.sql @@ -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; diff --git a/tradein-mvp/backend/tests/scrapers/test_domclick_detail.py b/tradein-mvp/backend/tests/scrapers/test_domclick_detail.py index 6cf14c12..dad63111 100644 --- a/tradein-mvp/backend/tests/scrapers/test_domclick_detail.py +++ b/tradein-mvp/backend/tests/scrapers/test_domclick_detail.py @@ -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 → diff --git a/tradein-mvp/backend/tests/test_3253_domclick_house_fields.py b/tradein-mvp/backend/tests/test_3253_domclick_house_fields.py index 909477da..862c98e9 100644 --- a/tradein-mvp/backend/tests/test_3253_domclick_house_fields.py +++ b/tradein-mvp/backend/tests/test_3253_domclick_house_fields.py @@ -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: diff --git a/tradein-mvp/backend/tests/test_snapshot_writer.py b/tradein-mvp/backend/tests/test_snapshot_writer.py index 6cde5c51..bd567917 100644 --- a/tradein-mvp/backend/tests/test_snapshot_writer.py +++ b/tradein-mvp/backend/tests/test_snapshot_writer.py @@ -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) diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/offer_price_history.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/offer_price_history.py index 6468e5d6..8e8eb726 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/offer_price_history.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/offer_price_history.py @@ -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 diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/detail.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/detail.py index 9dc77161..030345da 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/detail.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/detail.py @@ -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(""" diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/domclick/detail.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/domclick/detail.py index 943a4a92..f1fe6f5c 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/domclick/detail.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/domclick/detail.py @@ -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)", diff --git a/tradein-mvp/scripts/sql/286_dryrun.sql b/tradein-mvp/scripts/sql/286_dryrun.sql new file mode 100644 index 00000000..66e11266 --- /dev/null +++ b/tradein-mvp/scripts/sql/286_dryrun.sql @@ -0,0 +1,252 @@ +-- 286_dryrun.sql — прогон миграции 286 на проде БЕЗ записи (#3376). +-- +-- ЗАЧЕМ. Миграция 286 удаляет строки из offer_price_history. Прежде чем этому +-- уехать на деплой, надо увидеть на БОЕВЫХ данных: сколько строк она заберёт по +-- источникам, сколько соседей пересчитает, сколько процентов обнулит — и что +-- ВТОРОЙ прогон подряд не делает уже ничего (идемпотентность). +-- +-- КАК ЗАПУСКАТЬ (ничего не остаётся, в конце ROLLBACK): +-- ssh poincare +-- docker exec -i tradein-postgres psql -U -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;