"""Deactivate stale listings that have not been seen in >TTL days. Первоначально создано для avito (#759). Теперь поддерживает любой источник (listing_source) с опциональной фильтрацией по listing_segment. Ключевые решения: - Cian/Yandex не поддерживают full-coverage sweep -> паушальный TTL сломает живой инвентарь. DECISION: для yandex/cian деактивировать ТОЛЬКО listing_segment='vtorichka', TTL=30. novostroyki (9659 активных первичных строк) не трогаем. ИЗВЕСТНЫЙ ПРОБЕЛ (ревью TTL-CAP круг 2, 2026-08-15): этот скоуп уже, чем множество реально протухших строк -- живой замер на проде даёт cian/novostroyki 9 483 активных строки старше 60 суток, ни одна из них не деактивируется НИ ОДНОЙ джобой (внутри скоупа cian/vtorichka и yandex/vtorichka таких строк 0). Потолок cap_mult (см. CAP_MULT ниже) этот пробел не закрывает и закрыть не может -- он сжимает пул ВНУТРИ скоупа джобы, а не расширяет сам скоуп. NULL-сегмент (тот же замер круга 2 давал cian/NULL 211, yandex/NULL 523 строки старше 60 суток) закрыт отдельно ниже (null_segment_only, миграция 266) -- novostroyki-часть пробела остаётся: расширение скоупа туда отдельная задача (нужно сперва выяснить, поддерживает ли cian/yandex full-coverage sweep для novostroyki СЕЙЧАС, иначе паушальный TTL повторит инцидент, ради которого этот DECISION и принят) и намеренно НЕ входит в TTL-CAP. - avito: все сегменты (segments=None), TTL=10 дней -- поведение без изменений. - NULL-сегмент (легаси-строки до миграции 011 + жертвы бага в ON CONFLICT -- upsert никогда не пишет listing_segment повторно, поэтому раз рождённая NULL-строка сама себя не чинит даже при живой ежедневной досдаче) деактивируется ОТДЕЛЬНОЙ джобой per source (null_segment_only=True, миграция 266): явный `listing_segment IS NULL` предикат, а не ANY(:segments) -- этот оператор NULL никогда не матчит. Гейт здоровья/пол переобхода для этой джобы выключены (min_confirmations=0, revisit_floor_quantile=0) -- население нерепрезентативно мало (единицы подтверждений в сутки против сотен-тысяч у обычного vtorichka-среза), калиброванный под vtorichka порог держал бы джобу вечно skipped_unhealthy. - Строки НЕ удаляются -- история нужна для бэктеста (#667). - #2674: деактивация в той же транзакции пишет снимок listings_snapshots со статусом 'stale' за текущую дату -- «мы N суток не видели». Жёсткое 'closed' (площадка ответила 404) пишет только avito_detail_backfill: смешивать факт с догадкой дорого. Задача синхронная (DB-only, никаких внешних HTTP-вызовов) -- запускается kit-scheduler'ом через product_handlers._job_deactivate_stale (wildcard-handler deactivate_stale_*), по образцу snapshot_listing_sources / recompute_asking_to_sold_ratios (sync task в run_in_executor). TTL для avito берётся из settings.avito_stale_ttl_days (env AVITO_STALE_TTL_DAYS, default 10). """ from __future__ import annotations import logging from math import ceil from typing import Any from sqlalchemy import text from sqlalchemy.orm import Session from app.core.config import settings from app.services import scrape_runs as runs_mod logger = logging.getLogger(__name__) # Разрешённые колонки-таймстемпы для проверки свежести (staleness). # Whitelist ЖЁСТКО ограничивает f-string-интерполяцию имени колонки в SQL: # значение колонки НИКОГДА не берётся из пользовательского ввода без этой проверки. # - last_seen_at: дефолт; двигается при каждом «увиденном» объявлении в прогоне. # - scraped_at: честная свежесть для источников с нетрекаемым bulk-touch # last_seen_at (domklik: #2204 — bulk UPDATE двигает last_seen_at всем строкам # одним timestamp, поэтому TTL по last_seen_at бесполезен; scraped_at двигает # только реальный скрейп). _ALLOWED_STALENESS_COLUMNS = frozenset({"last_seen_at", "scraped_at"}) # ── Снимок «протухло» в дневной истории (#2674) ─────────────────────────────── # listings_snapshots.status до этого фикса был константой 'active' у всех строк # (394 299 на момент находки) — оба места вызова upsert_listing_snapshot передавали # литерал 'active', и это честно: там объявление ДЕЙСТВИТЕЛЬНО видели. А деактивация # TTL-задачей не оставляла в истории вообще никакого следа. Из-за этого дата снятия # объявления (лучший доступный сигнал «скорее всего продано») не запрашивалась из # истории, а восстанавливалась на глаз: последний показ + предполагаемый срок жизни. # # ПОЧЕМУ 'stale', А НЕ 'closed'. Эта задача НЕ знает, что объявление снято, — она # знает только, что МЫ его N суток не видели, а это разные факты, когда TTL короче # простоя обхода. Замер: прогон по домклику 02.08 снял 6131 объявление за раз (TTL # 14 суток против 12 суток простоя обхода) — с общим статусом это были бы 6131 # фальшивая «дата продажи» одной датой. Продукт про цены, смешивать факт с догадкой # дорого. Поэтому: # 'closed' — только путь 404: площадка ответила «нет» (avito_detail_backfill); # 'stale' — этот путь: «мы N суток не смотрели». # Дата всё равно фиксируется, но читатель отличает одно от другого. Ограничения # CHECK на колонке нет (проверено на проде), миграция не нужна — только COMMENT. # # Пишем снимок в ТОЙ ЖЕ транзакции, что и UPDATE флага: деактивация без снимка (или # наоборот) невозможна по построению — один statement, data-modifying CTE. # 1:1 по строкам: `stale` возвращает уникальные listings.id (PK), каждая даёт ровно # одну затронутую строку listings_snapshots (INSERT либо DO UPDATE — оба считаются # в rowcount), поэтому rowcount statement'а по-прежнему равен числу деактивированных. # price_rub берём из listings (NOT NULL в схеме) — это последняя известная цена. # ON CONFLICT: если снимок за сегодня уже есть (объявление видели активным утром, # а вечером сработал TTL) — только переводим статус в 'stale', цену не переписываем. _STALE_SNAPSHOT_TAIL = """ INSERT INTO listings_snapshots (listing_id, snapshot_date, run_id, price_rub, status, observed_at) SELECT id, CURRENT_DATE, CAST(:run_id AS bigint), price_rub, 'stale', NOW() FROM stale ON CONFLICT (listing_id, snapshot_date) DO UPDATE SET status = 'stale', observed_at = EXCLUDED.observed_at """ # ── Гейт по здоровью сбора (#2659) ──────────────────────────────────────────── # TTL отвечает на вопрос «объявление сняли?», а меряет «мы его давно не видели». # Пока обход здоров, разница мала. Когда обход лёг — разница равна всему инвентарю. # # Замер на проде, из-за которого этот гейт существует. Авито 10.07-26.07.2026: # 17 суток подряд без единой успешно собранной страницы, TTL=10 снял за этот отрезок # 9 033 строки; 1 270 из них потом доказанно вернулись живыми (снимки # listing_source_snapshots + текущий last_seen_at) — сбор восстановился, и объявления # оказались на месте. То есть «мы не смогли зайти» было прочитано как «объявление снято». # # ПОЧЕМУ НЕ ПО СТАТУСУ ban. Соблазн взять scrape_runs.status='banned' — ловушка: # Яндекс 18.07-30.07 — 5 прогонов в сутки, ВСЕ 'done', НОЛЬ 'banned', total_seen=0 # 13 суток подряд (снято ~839 строк vtorichka); # Домклик 20.07-30.07 — то же самое, 11 суток 'done' с total_seen=0, а 02.08 TTL # снял 6 131 строку разом (см. комментарий про 'stale' выше). # Оба провала для ban-детектора невидимы. Поэтому здоровье меряем НЕ статусом прогона, # а результатом: сколько строк источник реально подтвердил свежими за последние сутки. # # МЕТРИКА: count(*) по той же колонке свежести, что и сам TTL (last_seen_at или # scraped_at) и по тому же срезу source+segment, что и UPDATE. Одна колонка на обе # стороны — гейт нельзя обмануть bulk-touch'ем, который не двигает scraped_at (#2204). # # ПОРОГ. Ряд «подтверждений за 3 суток» по дням (восстановлен из listing_source_snapshots): # avito здоровые сутки 3542..6079, провал 10.07-26.07 — 0..970 → порог 1500; # yandex vtorichka здоровые 897..2206, провал — 0 → порог 500; # cian vtorichka 748..4329, провала не было → порог 500; # domklik по scraped_at сейчас 62/3 суток (сбор фактически стоит) → порог 200. # Пороги живут в default_params расписания (миграция 219), здесь только страховка # на случай незасеянного расписания. Асимметрия цены ошибки намеренная: пропущенная # деактивация чинится следующим прогоном, ложная — только повторным сбором, которого # может не быть. Поэтому при сомнении — пропускаем прогон. # # ПОТОЛОК: окно 3 суток годится, пока свипы источника ходят не реже чем раз в 3 дня. # Источник с более редкой каденцией будет блокироваться всегда — тогда окно нужно # растить до каденции, а не понижать порог. _HEALTH_WINDOW_DAYS = 3 # Страховка для расписаний без явного min_confirmations в default_params: ловит # полный ноль и близкое к нулю, но НЕ ловит частичный провал вроде avito 936-970 — # для этого нужен посчитанный по источнику порог из миграции 219. DEFAULT_MIN_CONFIRMATIONS = 500 _CONFIRMATIONS_SEGMENT_FILTER = "\n AND listing_segment = ANY(CAST(:segments AS text[]))" # NULL-сегмент: `= ANY(...)` НИКОГДА не матчит NULL (SQL, не баг), поэтому # для null_segment_only-режима нужен отдельный явный предикат IS NULL, а не элемент # в :segments. См. _build_null_segment_sql ниже -- тот же принцип для самого UPDATE. _CONFIRMATIONS_NULL_SEGMENT_FILTER = "\n AND listing_segment IS NULL" # ── Пол TTL по измеренному циклу переобхода (#2659) ─────────────────────────── # Гейт выше отвечает на вопрос «источник вообще собирается?». Он НЕ отвечает на # вопрос, из-за которого заведён #2659: «а достаточно ли ttl_days, чтобы молчание # означало снятие?». Пока свип возвращается к строке реже, чем раз в ttl_days, # TTL меряет НАШУ выборку, а не жизнь объявления, — и источник при этом полностью # здоров, так что гейт молчит. # # ЗАМЕР НА ПРОДЕ 2026-08-09, из-за которого этот пол существует. # С момента деплоя гейта (06.08) TTL снял 1 028 строк; 127 из них (12.4%) УЖЕ снова # активны — свип нашёл их живыми через 1-3 суток и вернул сам (upsert в # scraper_kit/base.py ставит is_active = true). В единственном городе с настоящим # покрытием доля ложных снятий 100%: # cian Екатеринбург 103 снято → 103 снова активны # yandex Екатеринбург 24 снято → 24 снова активны # cian/yandex без города 901 снято → 0 вернулись (их свип не обходит вовсе) # Возраст на момент снятия у всех 127: 29.9..30.3 суток при TTL=30 — то есть TTL # срабатывал ровно на границе, а свип возвращался к строке на 31-34-е сутки. # # ПОЧЕМУ ЭТО НЕ ЛЕЧИТСЯ НОВОЙ КОНСТАНТОЙ. Разрывы переобхода, суток # (listing_source_snapshots, 40 суток, посчитано по срезу TTL-джобы): # источник/сегмент p90 p99 TTL сейчас TTL/p99 # domklik vtorichka 1.9 3.1 14 4.5 ← сплошное суточное покрытие # cian vtorichka 10.9 26.6 30 1.1 # yandex vtorichka 5.7 43.0 30 0.7 # avito vtorichka 29.1 42.1 10 0.24 ← отсюда 9 033 строки # Домклик — контрольная группа: при почти полном суточном обходе TTL=14 лежит в # 4.5 раза выше хвоста, и снятие у него действительно означает снятие. У остальных # трёх порог ниже собственного хвоста обхода — руками подобранное число и есть # корень #2659, поэтому чинить его вторым руками подобранным числом бессмысленно. # # ЧТО МЕРЯЕМ ВМЕСТО КОНСТАНТЫ: факт, а не оценку. «Какой самый большой возраст, при # котором свип за последнее окно ДОКАЗАЛ, что объявление живо» — то есть насколько # старую строку он только что нашёл на площадке. Если свип буквально вчера вернул к # жизни строку, молчавшую 40 суток, то 30 суток молчания не доказывают ничего. # Пол = квантиль этого распределения, эффективный TTL = max(ttl_days, пол). # # Считается по ТОМУ ЖЕ срезу (source + segments) и по ТОЙ ЖЕ колонке свежести, что # и UPDATE. Предыдущее наблюдение берётся из listing_source_snapshots — единственной # истории свежести, что у нас есть; расхождение listings. и # listing_sources.last_seen_at замерено на проде и не превышает 0.5 суток в среднем # (максимум 0), что на шкале 30-70 суток шум. # # КВАНТИЛЬ — калибровочная ручка, не догма. 0.99 подобран по требованию «пол обязан # накрыть 127 доказанных ложных снятий», у которых возраст был 29.9..30.3: замер # того же запроса на проде даёт 34.0 для cian/vtorichka и 74.3 для yandex/vtorichka. # Ниже 0.99 опускать нельзя без нового замера. Ручка живёт в default_params # расписания (revisit_floor_quantile), 0 -> пол выключен. # # ПОБОЧНЫЙ ЭФФЕКТ, КОТОРЫЙ ЗДЕСЬ НАМЕРЕННЫЙ: после провала сбора хвост разрывов # распухает (свип разгребает завал и находит очень старые строки), пол поднимается, # и деактивация замирает сама — без отдельного детектора банов. Когда завал разобран, # хвост схлопывается и пол опускается обратно. Это ровно то поведение, которого # issue просил от «гейта по банам», но выраженное через результат, а не через причину. # # ПОТОЛОК: пол не может превысить глубину истории снимков. Если снимок за нужную # дату не писался (дыры на проде есть — 30.07, 01.08), берётся ближайший более # ранний; при полном отсутствии снимков пол не считается и TTL остаётся как задан. DEFAULT_REVISIT_FLOOR_QUANTILE = 0.99 _REVISIT_FLOOR_SEGMENT_FILTER = "\n AND l.listing_segment = ANY(CAST(:segments AS text[]))" _REVISIT_FLOOR_NULL_SEGMENT_FILTER = "\n AND l.listing_segment IS NULL" # ── Потолок эффективного TTL (положительная обратная связь пола, найдено 2026-08-15) ── # У пола выше нет верхней границы: max(ttl_days, пол) может расти неограниченно. # ЗАМЕР НА ПРОДЕ (уточнён 2026-08-15 после разбора): у yandex counters держали # ttl_days_effective 75/75/75/39/52/54 шесть прогонов подряд при deactivated=0 — # пол реально разгонялся без верхней границы, и потолок закрывает именно это. # ЧЕГО ПОТОЛОК НЕ ДЕЛАЕТ: он НЕ сжимает пул «активных». Замер показал 0 # деактивируемых строк на всех четырёх джобах и до, и после калибровки. Цифра # «23 687 из 44 744 не подтверждались >7 суток» относится ко ВСЕМ источникам # сразу, и две трети её — новостройки, которых оценщик не берёт. У avito # просроченных ноль. Раздутый пул, влияющий на оценку, лежит в строках с ПУСТЫМ # сегментом и чинится отдельной джобой, не этим потолком. # # МЕХАНИЗМ ПЕТЛИ: медленный обход поднимает пол (он же квантиль разрывов переобхода) # -> высокий пол продлевает жизнь снятым лотам дольше, чем к ним успевает вернуться # свежий обход -> пул «активных» раздувается «протухшими» строками -> следующий замер # пола на том же раздутом пуле оказывается ещё выше. Без верхней границы это не # самокорректирующийся пол, а положительная обратная связь. # # CAP_MULT = 2 -- эффективный TTL не может превысить удвоенный заданный оператором # ttl_days. Пол по-прежнему может его поднять (ради #2659 -- см. комментарий выше: # ложные снятия при неполном покрытии обхода), но не бесконечно. Почему именно 2, а # не 3 или 1.5: вдвое — это ещё «мы искренне не уверены, что молчание значит # снятие», не «источник вообще умер». Дальнейший рост пола сигнализирует не о # медленном, но живом обходе, а о мёртвом источнике -- для ЭТОГО случая уже есть # отдельный гейт по здоровью (min_confirmations) выше в этой же функции, который # выключает деактивацию целиком, а не растягивает TTL до бесконечности. Калибровочная # ручка, не догма -- при новом замере можно пересмотреть, как и revisit_floor_quantile. # # ПОЧЕМУ MULT, А НЕ ФИКСИРОВАННОЕ ЧИСЛО СУТОК -- И ГДЕ ЭТА ФОРМА ЛОМАЕТСЯ. Множитель # от ttl_days даёт разный АБСОЛЮТНЫЙ потолок на разных источниках: cian/yandex # (ttl=30) -> 60 суток, avito (ttl=10) -> 20 суток, domklik (ttl=14) -> 28 суток. Это # ломается ровно там, где абсолютный хвост переобхода источника НЕ пропорционален его # ttl_days. Замер (_REVISIT_TAIL, 40 суток): avito p99 = 42.1 сут -- ВЫШЕ его же # потолка 20. То есть для avito дефолтный CAP_MULT=2 может резать ttl ниже # собственного хвоста обхода -- ровно тот false-kill, ради которого пол вообще # заведён (см. комментарий выше). domklik (потолок 28 при хвосте 3.1) разрыва не # имеет -- множитель 2 для него калиброван верно. # # YANDEX -- ТА ЖЕ ДЫРА, НАЙДЕНА ПОЗЖЕ (ревью круга 3, 2026-08-15). Строка выше до # этой правки утверждала, что cian/yandex с потолком 60 тоже в порядке -- это было # верно для cian (live-пол сейчас 31.1), но НЕ для yandex: ЖИВЫЕ полы из # scrape_runs.counters (deactivate_stale_yandex, 2026-08-10..08-15) -- 75/75/75/39/ # 52/54, а прямой live-замер той же percentile_disc(0.99)-формулы сегодня даёт 79.2 # (n=1961 подтверждений за 3 суток). И то, и другое ВЫШЕ потолка 60 -- тот же # false-kill класс, что у avito, статический p99=43.0 (_REVISIT_TAIL) для yandex # устарел и вводит в заблуждение. cap_mult для yandex откалиброван отдельной # миграцией (265_deactivate_stale_yandex_cap_mult.sql, cap_mult=3 -> потолок 90) -- # см. её комментарий про то, почему это НЕ меняет число деактивированных строк # следующим прогоном (0 активных строк источника старше 39 суток на момент замера). # # ПОЭТОМУ cap_mult -- параметр функции (как revisit_floor_quantile, min_confirmations), # не голая константа: default = CAP_MULT для источников, где 2x достаточно (cian, # domklik), но расписание может переопределить через default_params (JSON-колонка # scrape_schedules, ключ "cap_mult") для источника с непропорционально длинным # хвостом -- см. миграции для avito (cap_mult=6, потолок 60, с запасом выше # статического p99=42.1 и живого прод-пика 52, замеренного 2026-08-10..12) и yandex # (cap_mult=3, потолок 90, с запасом выше живого пола 79.2, замеренного 2026-08-15). CAP_MULT = 2 # ── Гейт деградации пола (PR-B, #2659 продолжение) ───────────────────────────── # У джобы деактивации до сих пор не было НИ ОДНОГО ограничителя ОБЪЁМА снятия: # min_confirmations может пропустить прогон целиком, пол поднимает TTL, cap_mult # ограничивает сам пол -- но ни один из них не смотрит на РАЗМЕР ВЫБОРКИ, из # которой percentile_disc посчитал квантиль (floor_n_pairs, PR-A). Это другая # величина, чем min_confirmations: confirmations считает ВСЕ строки, увиденные # свежими за окно, floor_n_pairs -- только те из них, чья свежесть СДВИНУЛАСЬ # относительно предыдущего снимка (см. _revisit_floor_from_where_sql). Источник # может быть формально здоров (confirmations высокий) при вырожденном floor_n_pairs # -- ровно то, что уже наблюдалось при переходе на change-only модель снимков # (PR-A: equality-join терял 95.6% пар при здоровом источнике). # # ПОЧЕМУ ЭТОТ ГЕЙТ САМОКАЛИБРУЮЩИЙСЯ, А НЕ РУЧНОЙ ПОРОГ ПО ИСТОЧНИКУ. В отличие # от min_confirmations (числа калибровались по восстановленному ряду за конкретные # месяцы, миграция 219) у floor_n_pairs как счётчика истории почти нет -- он # появился только в PR-A. Поэтому главный сигнал -- ОТНОСИТЕЛЬНЫЙ: падение # floor_n_pairs относительно ПРЕДЫДУЩЕГО УСПЕШНОГО прогона ТОГО ЖЕ расписания # (сравнение "прогон сам с собой", свежей калибровки по чужим числам не требует). # Предыдущее значение читается из scrape_runs.counters по scrape_runs.source # ТЕКУЩЕГО прогона (см. _PREVIOUS_FLOOR_N_PAIRS_SQL ниже) -- это имя РАСПИСАНИЯ # (deactivate_stale_avito / _yandex / _cian / _yandex_null_segment / …, миграции # 090/115/266), НЕ listing_source из параметров функции: у cian и его же # null-сегмент джобы listing_source один и тот же ('cian'), но это две независимые # серии прогонов с несопоставимым масштабом выборки (vtorichka -- сотни-тысячи # пар, null-сегмент -- единицы, см. модульный докстринг про 266), сравнивать их # друг с другом было бы категориальной ошибкой. # # АБСОЛЮТНЫЙ min_floor_pairs -- страховка на случай, когда истории предыдущего # прогона ещё нет (первый прогон после деплоя гейта) либо она сама уже была # вырожденной (иначе относительный порог сравнивал бы вырожденное с вырожденным # и молчал бы вечно). DEFAULT_MIN_FLOOR_PAIRS = 30 -- ниже этого n percentile_disc # при квантиле 0.99 неотличим от максимума выборки (одно случайное значение решает # результат), это статистический минимум устойчивости хвостового квантиля, а НЕ # число, снятое с прод-ряда floor_n_pairs -- собственной истории по нему пока нет. # 0 -> абсолютная проверка отключена (тот же идиом, что у min_confirmations и # revisit_floor_quantile: 0 = гейт снят); относительный дрейф-гейт при этом # продолжает работать самостоятельно. DEFAULT_MIN_FLOOR_PAIRS = 30 # DEFAULT_FLOOR_DROP_RATIO = 5 -- floor_n_pairs упал МИНИМУМ впятеро относительно # предыдущего успешного прогона того же расписания. Собственной калибровки по # floor_n_pairs пока нет (счётчик из PR-A) -- порядок величины заимствован у # соседнего, уже проверенного на реальном инциденте гейта: здоровый разброс avito # по суткам confirmations 2542..6079 (_AVITO_HEALTHY_DAYS, # tests/test_deactivate_stale_health_gate.py) -- это ~2.4x, то есть штатный # день-в-день шум заведомо НЕ достигает 5x. Ratio=5 оставляет двукратный запас # над этим шумом и всё ещё ловит провал порядка 10.07-26.07 (падение в 5-10+ раз). # floor_n_pairs -- другая метрика (подмножество confirmations, только реально # переобойдённые строки), но природа шума та же (день-в-день колебание объёма # обхода), поэтому заимствование порядка величины -- разумная отправная точка, # а не число с потолка; пересмотреть, когда floor_n_pairs накопит собственную # историю на проде. DEFAULT_FLOOR_DROP_RATIO = 5.0 def _revisit_floor_from_where_sql(staleness_column: str, segment_filter: str) -> str: """FROM..WHERE, общий для _build_revisit_floor_sql и _build_revisit_floor_pairs_count_sql. Вынесено в отдельную функцию, а не продублировано, ЦЕЛЕНАПРАВЛЕННО: пункт 3 PR-A (#2659 продолжение) требует, чтобы n_pairs в counters считался «по тому же срезу, что даёт выборку для percentile_disc» -- дословно, а не приблизительно. Общий текст гарантирует это по построению; раздельные копии SELECT count(*) и SELECT percentile_disc(...) рано или поздно разошлись бы независимой правкой одной из двух. JOIN LATERAL — предшественник ищется ПО СТРОКЕ (per listing_source_id), не по глобальной дате: `ORDER BY s.snapshot_date DESC LIMIT 1` берёт последний снимок ЭТОГО listing_source_id не позже якоря (CURRENT_DATE - health_window_days), независимо от того, писалась ли строка именно на дату якоря. Подробно — почему именно так и что было раньше — см. докстринг _build_revisit_floor_sql. """ return f""" FROM listings l JOIN listing_sources ls ON ls.listing_id = l.id AND ls.ext_source = l.source JOIN LATERAL ( SELECT s.last_seen_at FROM listing_source_snapshots s WHERE s.listing_source_id = ls.id AND s.snapshot_date <= CURRENT_DATE - CAST(:health_window_days AS integer) ORDER BY s.snapshot_date DESC LIMIT 1 ) prev ON true WHERE l.source = :listing_source AND l.{staleness_column} > NOW() - CAST(:health_window_days || ' days' AS interval) AND l.{staleness_column} > prev.last_seen_at{segment_filter} """ def _build_revisit_floor_sql( staleness_column: str, *, with_segments: bool, null_segment_only: bool = False ) -> Any: """Квантиль возраста, при котором свип за окно ДОКАЗАЛ, что строка жива. Пара «предыдущее наблюдение (снимок) → текущее наблюдение (listings)» даёт разрыв переобхода в сутках; берём его квантиль по срезу source+segments. Только строки, у которых свежесть реально сдвинулась, — то есть выжившие, а не «мы к ним не приходили». «ПРЕДЫДУЩЕЕ НАБЛЮДЕНИЕ» -- ЧТО ЭТО ТОЧНО (PR-A, #2659 продолжение). Это ПОСЛЕДНЯЯ строка listing_source_snapshots для данного listing_source_id, чей snapshot_date не позже якоря (CURRENT_DATE - health_window_days) -- `ORDER BY snapshot_date DESC LIMIT 1` в LATERAL-подзапросе _revisit_floor_from_where_sql. РАВЕНСТВО ПО ДАТЕ (`prev.snapshot_date = <якорь>`) ЗДЕСЬ ЗАПРЕЩЕНО, и вот почему: 1) listing_source_snapshots переходит на модель «строка на изменение» (пишется только когда значение отличается от предыдущего снимка, не ежедневно). В этой модели equality-join по дате теряет подавляющее большинство пар: на 2026-08-20 «изменениями» являются 4 458 строк из 101 795 -- 95.6% пар для percentile_disc пропадают, выборка схлопывается до n=3-4, а percentile_disc(0.99) на такой выборке вырождается в максимум из трёх чисел, а не в реальный квантиль хвоста. 2) Дыры в ЕЖЕДНЕВНОЙ истории уже ломали equality-join и ДО перехода на change-only модель -- на проде отмечены разрывы 03-14.06 (12 суток подряд), 03-04.07, 12.07, 26.07, 30-31.07, 01.08. В эти дни equality-join давал n_pairs=0, floor_days=NULL, и пол молча не считался -- TTL оставался как задан, хотя история для «ближайшего более раннего» снимка (см. комментарий про потолок выше в модульном докстринге) в базе была, просто не РОВНО на эту дату. Семантика при этом не меняется: между двумя изменениями last_seen_at по определению постоянен (иначе строка была бы новым изменением), поэтому «последний снимок не позже якоря» и «снимок ровно на дату якоря» при СПЛОШНОЙ ежедневной истории дают одно и то же число -- разница проявляется только там, где equality-join был неверен и раньше (гэпы) либо станет неверен при переходе на change-only (почти повсеместно). null_segment_only=True переопределяет with_segments -- IS NULL вместо ANY(:segments) (ANY никогда не матчит NULL). На практике для null_segment_only-джобы этот запрос не строится вовсе (revisit_floor_quantile=0 -- см. модульный докстринг), но вариант нужен для корректности, если порог когда-нибудь включат. staleness_column уже прошёл whitelist-проверку в deactivate_stale_listings. Значения — param-binding, psycopg v3 safe (CAST(... AS ...), никаких :param::type). """ if null_segment_only: segment_filter = _REVISIT_FLOOR_NULL_SEGMENT_FILTER elif with_segments: segment_filter = _REVISIT_FLOOR_SEGMENT_FILTER else: segment_filter = "" return text( f""" SELECT percentile_disc(CAST(:revisit_quantile AS double precision)) WITHIN GROUP ( ORDER BY EXTRACT(epoch FROM (l.{staleness_column} - prev.last_seen_at)) / 86400.0 ) {_revisit_floor_from_where_sql(staleness_column, segment_filter)} """ ) def _build_revisit_floor_pairs_count_sql( staleness_column: str, *, with_segments: bool, null_segment_only: bool = False ) -> Any: """count(*) пар, из которых percentile_disc в _build_revisit_floor_sql берёт квантиль. Наблюдательность (PR-A, #2659 продолжение) -- НИЧЕГО не блокирует сейчас (гейт по деградации выборки, если он когда-нибудь понадобится, -- отдельная задача, PR-B). Пишется в counters["floor_n_pairs"] ДО того, как схлопнувшаяся выборка станет видна только по повторению прод-инцидента, ради которого весь пол заведён (см. ЗАМЕР НА ПРОДЕ 2026-08-09 в модульном докстринге). Тот же срез, что и percentile_disc -- ОБЩАЯ функция _revisit_floor_from_where_sql, не копия WHERE, см. её докстринг про то, почему это важно. """ if null_segment_only: segment_filter = _REVISIT_FLOOR_NULL_SEGMENT_FILTER elif with_segments: segment_filter = _REVISIT_FLOOR_SEGMENT_FILTER else: segment_filter = "" return text( f""" SELECT count(*) {_revisit_floor_from_where_sql(staleness_column, segment_filter)} """ ) def _build_confirmations_sql( staleness_column: str, *, with_segments: bool, null_segment_only: bool = False ) -> Any: """SELECT count(*) подтверждённых за окно строк — тот же срез, что и у UPDATE. null_segment_only=True переопределяет with_segments -- IS NULL вместо ANY(:segments). Для null_segment_only-джобы min_confirmations=0 по умолчанию (см. модульный докстринг), так что на практике этот путь не строится -- оставлен для корректности. staleness_column уже прошёл whitelist-проверку в deactivate_stale_listings. Значения (:listing_source, :health_window_days, :segments) — param-binding, psycopg v3 safe (CAST(... AS ...), никаких :param::type). """ if null_segment_only: segment_filter = _CONFIRMATIONS_NULL_SEGMENT_FILTER elif with_segments: segment_filter = _CONFIRMATIONS_SEGMENT_FILTER else: segment_filter = "" return text( f""" SELECT count(*) FROM listings WHERE source = :listing_source AND {staleness_column} > NOW() - CAST(:health_window_days || ' days' AS interval){segment_filter} """ ) def _build_all_segments_sql(staleness_column: str) -> Any: """UPDATE без фильтра по сегменту: все сегменты для данного source. staleness_column уже прошёл whitelist-проверку в deactivate_stale_listings, поэтому f-string-подстановка имени колонки безопасна. Значения (:listing_source, :ttl_days, :run_id) остаются param-binding — psycopg v3 safe (никаких :param::type). """ return text( f""" WITH stale AS ( UPDATE listings SET is_active = false WHERE source = :listing_source AND is_active = true AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval) RETURNING id, price_rub ) {_STALE_SNAPSHOT_TAIL} """ ) def _build_segments_sql(staleness_column: str) -> Any: """UPDATE с фильтром по сегменту (segments задан): только указанные сегменты. staleness_column уже прошёл whitelist-проверку. = ANY(CAST(:segments AS text[])) -- psycopg v3 адаптирует Python list -> text[]. """ return text( f""" WITH stale AS ( UPDATE listings SET is_active = false WHERE source = :listing_source AND is_active = true AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval) AND listing_segment = ANY(CAST(:segments AS text[])) RETURNING id, price_rub ) {_STALE_SNAPSHOT_TAIL} """ ) def _build_null_segment_sql(staleness_column: str) -> Any: """UPDATE строго по listing_segment IS NULL (null_segment_only=True). НЕ переиспользует _build_segments_sql: `= ANY(CAST(:segments AS text[]))` никогда не матчит NULL (SQL-семантика, не баг -- та же ловушка задокументирована выше у novostroyki-гарда), поэтому NULL-сегмент не выразить через список segments и нужен отдельный явный предикат. Целенаправленно НЕ трогает 'vtorichka'/'novostroyki' -- их деактивация идёт через _build_segments_sql в отдельных, уже существующих джобах. staleness_column уже прошёл whitelist-проверку. Без :segments-параметра вовсе. """ return text( f""" WITH stale AS ( UPDATE listings SET is_active = false WHERE source = :listing_source AND is_active = true AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval) AND listing_segment IS NULL RETURNING id, price_rub ) {_STALE_SNAPSHOT_TAIL} """ ) # Дефолтные (last_seen_at) варианты SQL -- сохранены как модульные константы для # обратной совместимости (тесты читают .text, product_handlers/scheduler не менялись). _DEACTIVATE_SQL_ALL_SEGMENTS = _build_all_segments_sql("last_seen_at") _DEACTIVATE_SQL_SEGMENTS = _build_segments_sql("last_seen_at") # Алиас для обратной совместимости -- использовался в тестах через task_mod._DEACTIVATE_SQL _DEACTIVATE_SQL = _DEACTIVATE_SQL_ALL_SEGMENTS # ── Гейт деградации пола: чтение предыдущего успешного прогона (PR-B) ────────── # source сравнивается по scrape_runs.source ТЕКУЩЕГО прогона (подзапрос по # :run_id), а НЕ по listing_source -- см. комментарий у DEFAULT_MIN_FLOOR_PAIRS # про то, почему это разные вещи. status='done' + id != :run_id исключают сам # текущий прогон (на момент этого запроса он ещё 'running', так что status='done' # уже достаточно, id != добавлен как явная защита от совпадения). counters ->> # 'floor_n_pairs' IS NOT NULL заменяет `?`-оператор существования ключа -- # семантически то же самое (NULL, если ключа нет), без сомнений по поводу # взаимодействия `?` с bind-параметрами psycopg v3 в этом же тексте. _PREVIOUS_FLOOR_N_PAIRS_SQL = text( """ SELECT CAST(prev.counters ->> 'floor_n_pairs' AS integer) AS floor_n_pairs FROM scrape_runs prev WHERE prev.source = (SELECT source FROM scrape_runs WHERE id = :run_id) AND prev.status = 'done' AND prev.id != :run_id AND prev.counters ->> 'floor_n_pairs' IS NOT NULL ORDER BY prev.id DESC LIMIT 1 """ ) # ── Потолок объёма снятия -- аварийный, НЕ рабочий (PR-B, #2659 продолжение) ─── # У джобы деактивации никогда не было ограничителя ОБЪЁМА снятия: UPDATE идёт # одним statement'ом без LIMIT, counters["deactivated"] = result.rowcount -- при # обвале любого из гейтов выше снимается сколько снимется. # # ПОЧЕМУ НЕТ ПОРОГА ПО ДОЛЕ ПУЛА. Естественный кандидат -- "не больше N% активного # пула источника за один прогон" -- но пул, с которым эта джоба реально работает, # НИКОГДА не записывался: восстановить его постфактум можно только по # listing_sources, а deactivated считает строки listings -- разная гранулярность. # Проверено на историческом ряду снятий: доля деактивированного от восстановленного # пула -- 40%, 97%, 317%, 1652% -- числа не образуют осмысленного ряда, калибровать # порог не на чем. Поэтому здесь метрика ПОКА ТОЛЬКО ПИШЕТСЯ -- deactivation_candidates # (preflight count(*) ПО ТОМУ ЖЕ предикату, что исполнит UPDATE), active_pool # (все активные строки этого source, без фильтра по сегменту -- см. докстринг # deactivate_stale_listings, Returns), deactivated_pct. Долевой порог будет # выставлен ПОЗЖЕ, когда deactivated_pct накопит собственную историю на живых # прогонах именно этой джобы (а не на реконструкции задним числом). # # АВАРИЙНЫЙ ПОРОГ ЕСТЬ -- max_deactivated, абсолютное число, блокирует ДО UPDATE. # Исторический максимум ЛЕГИТИМНОГО снятия (owner-подтверждено): avito 06.06.2026 # -- 9 300 строк, затем 6 531 / 6 131 / 4 959 / 3 909 (последнее owner явно # подтвердил как здоровую чистку). DEFAULT_MAX_DEACTIVATED = 15 000 заведомо НЕ # отбивает ни один из этих легитимных прогонов (запас 1.6x над историческим # максимумом), но ловит катастрофу на порядок крупнее любого известного здорового # снятия -- это НЕ откалиброванный рабочий порог (как min_confirmations или # floor_drop_ratio выше), а предохранитель последней инстанции. DEFAULT_MAX_DEACTIVATED = 15000 _ACTIVE_POOL_SQL = text( """ SELECT count(*) FROM listings WHERE source = :listing_source AND is_active = true """ ) def _build_all_segments_candidates_count_sql(staleness_column: str) -> Any: """count(*) кандидатов на деактивацию -- ТОТ ЖЕ предикат, что WHERE в _build_all_segments_sql (см. её докстринг). Preflight ДО UPDATE (PR-B): даёт counters["deactivation_candidates"] и питает аварийный потолок max_deactivated -- при abort ни один UPDATE ещё не исполнялся. Предикат продублирован текстуально, а не вынесен в общую функцию с _build_all_segments_sql: рефакторинг уже протестированных UPDATE-builder'ов вне скоупа PR-B. Синхронность с UPDATE закреплена тестом test_candidates_predicate_matches_update_predicate (tests/test_deactivate_stale_deactivation_cap.py). """ return text( f""" SELECT count(*) FROM listings WHERE source = :listing_source AND is_active = true AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval) """ ) def _build_segments_candidates_count_sql(staleness_column: str) -> Any: """count(*) кандидатов -- ТОТ ЖЕ предикат, что WHERE в _build_segments_sql. См. докстринг _build_all_segments_candidates_count_sql выше. """ return text( f""" SELECT count(*) FROM listings WHERE source = :listing_source AND is_active = true AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval) AND listing_segment = ANY(CAST(:segments AS text[])) """ ) def _build_null_segment_candidates_count_sql(staleness_column: str) -> Any: """count(*) кандидатов -- ТОТ ЖЕ предикат, что WHERE в _build_null_segment_sql. См. докстринг _build_all_segments_candidates_count_sql выше. """ return text( f""" SELECT count(*) FROM listings WHERE source = :listing_source AND is_active = true AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval) AND listing_segment IS NULL """ ) def deactivate_stale_listings( db: Session, run_id: int, *, listing_source: str, ttl_days: int, segments: list[str] | None = None, staleness_column: str = "last_seen_at", min_confirmations: int = 0, health_window_days: int = _HEALTH_WINDOW_DAYS, revisit_floor_quantile: float = 0.0, null_segment_only: bool = False, cap_mult: float = CAP_MULT, min_floor_pairs: int = DEFAULT_MIN_FLOOR_PAIRS, floor_drop_ratio: float = DEFAULT_FLOOR_DROP_RATIO, max_deactivated: int = DEFAULT_MAX_DEACTIVATED, ) -> dict[str, int]: """Пометить is_active=false объявления, чья свежесть старше ttl_days дней. Параметры: listing_source: значение поля source в таблице listings ('avito', 'yandex', 'cian', 'domklik'). ttl_days: количество дней TTL; объявления старше этого порога деактивируются. segments: если задан -- деактивировать только объявления с указанными listing_segment значениями. None -> все сегменты (поведение avito по умолчанию). Несовместимо с null_segment_only=True (см. ниже). staleness_column: колонка-таймстемп, по которой считается свежесть. Whitelist {"last_seen_at", "scraped_at"} — иначе ValueError ДО любого SQL. Дефолт last_seen_at. Для domklik (#2204) — scraped_at: нетрекаемый bulk-touch двигает last_seen_at всем строкам одним timestamp, поэтому честная свежесть = scraped_at (двигается только реальным скрейпом). min_confirmations: гейт по здоровью сбора (#2659). Сколько строк источник должен был подтвердить свежими за health_window_days суток, чтобы деактивации вообще разрешалось исполниться. 0 -> гейт выключен (так вызывают старые тесты и совместимая обёртка); реальные значения приходят из default_params расписания, см. миграцию 219 и комментарий выше. health_window_days: окно подтверждений для гейта, суток. Дефолт 3. revisit_floor_quantile: пол TTL по измеренному циклу переобхода (#2659). Квантиль возраста, при котором свип за окно ДОКАЗАЛ строку живой; эффективный TTL = min(max(ttl_days, этот пол), ttl_days * cap_mult) -- пол поднимает TTL, но не выше потолка. 0 -> пол выключен (так вызывают старые тесты и совместимая обёртка), рабочее значение — DEFAULT_REVISIT_FLOOR_QUANTILE, см. комментарий выше. null_segment_only: True -> WHERE фильтрует `listing_segment IS NULL` вместо ANY(:segments). Требует segments=None (иначе ValueError -- смешивать бессмысленно, это два непересекающихся среза). Для этого среза гейт/пол обычно держат выключенными (min_confirmations=0, revisit_floor_quantile=0, см. миграцию 266 и модульный докстринг) -- население слишком мало для откалиброванных под полноценный vtorichka-свип порогов. cap_mult: множитель потолка эффективного TTL (см. комментарий у модульной константы CAP_MULT). Дефолт -- сама CAP_MULT=2, но параметр, а НЕ голая константа: источник с непропорционально длинным хвостом переобхода относительно своего ttl_days (avito: p99=42.1 при ttl=10 -> дефолтный потолок 20 режет ниже хвоста) может переопределить его через default_params расписания (ключ "cap_mult"), не трогая остальные источники. Итоговый потолок = ttl_days * cap_mult. Применяется и к null_segment_only-джобе, но там гейт/пол выключены (см. выше), так что на практике не участвует. min_floor_pairs: абсолютный порог гейта деградации пола (PR-B, #2659 продолжение). Ниже этого числа floor_n_pairs -- прогон пропускается (skipped_floor_degraded), НЕЗАВИСИМО от того, есть ли предыдущий прогон для сравнения. Дефолт DEFAULT_MIN_FLOOR_PAIRS=30, 0 -> отключает только эту (абсолютную) часть гейта -- см. комментарий у константы. Участвует ТОЛЬКО когда revisit_floor_quantile > 0 (гейт живёт внутри того же блока, что и сам пол -- вырожденную выборку нечем измерить, если пол вообще не считается). floor_drop_ratio: относительный порог того же гейта. floor_n_pairs упал СТРОГО более чем в floor_drop_ratio раз относительно предыдущего УСПЕШНОГО прогона того же расписания (scrape_runs.source, см. _PREVIOUS_FLOOR_N_PAIRS_SQL) -- прогон пропускается. Нет предыдущего прогона (первый после деплоя / история ещё не накопилась) -> эта часть гейта молчит, работает только min_floor_pairs. Дефолт DEFAULT_FLOOR_DROP_RATIO=5.0. max_deactivated: аварийный (НЕ рабочий, см. комментарий у константы) потолок объёма снятия за один прогон. Preflight count(*) кандидатов ПО ТОМУ ЖЕ предикату, что и сам UPDATE -- если кандидатов больше max_deactivated, прогон abort'ится ДО UPDATE. Дефолт DEFAULT_MAX_DEACTIVATED=15000. Участвует ТОЛЬКО когда revisit_floor_quantile > 0 -- на проде это верно для всех активных расписаний деактивации, кроме null_segment_only-джоб (миграция 266), чей пул на 2-3 порядка меньше порога и потолок для них физически не может сработать (см. комментарий у DEFAULT_MAX_DEACTIVATED). Sync (вызывается scheduler-триггером в executor, как snapshot_listing_sources). Один statement в транзакции: UPDATE флага + снимок 'stale' в listings_snapshots (data-modifying CTE, #2674). Финализирует scrape_runs (mark_done / mark_failed). Returns {"deactivated": N} -- количество обновлённых строк (1:1 со снимками). Если гейт здоровья не пропустил прогон: {"deactivated": 0, "confirmations": N, "skipped_unhealthy": 1} и НИ ОДНА строка не тронута. Если пол переобхода включён (revisit_floor_quantile > 0): дополнительно {"floor_n_pairs": N} -- размер выборки, из которой percentile_disc посчитал квантиль (PR-A), и (если найден предыдущий успешный прогон того же расписания) {"floor_n_pairs_previous": N}. Если гейт деградации пола не пропустил прогон (PR-B): {"skipped_floor_degraded": 1} и НИ ОДНА строка не тронута -- ни UPDATE, ни preflight-потолок ниже не исполнялись. Если пол при этом реально поднял TTL: дополнительно {"revisit_floor_days": N, "ttl_days_effective": N}. Если пол упёрся в потолок cap_mult: дополнительно {"ttl_floor_capped": 1, "ttl_days_floor_raw": N} -- N это то, во что пол поднял бы TTL БЕЗ потолка. Если revisit_floor_quantile > 0 и прогон не был abort'нут ни одним из гейтов выше -- PR-B добавляет наблюдательные {"deactivation_candidates": N, "active_pool": N, "deactivated_pct": N} (preflight ДО UPDATE, тот же предикат, что и сам UPDATE). Если кандидатов больше max_deactivated: {"skipped_cap_exceeded": 1} и НИ ОДНА строка не тронута. Raises: ValueError: если staleness_column не входит в whitelist, ИЛИ ttl_days <= 0 (проверка ДО SQL, никакой интерполяции пользовательского ввода в запрос; ttl_days<=0 в WHERE-условии last_seen_at < NOW() - INTERVAL 'N days' матчит практически весь активный пул -- без явного guard'а потолок (ttl_days * cap_mult <= 0) к тому же перебивал бы пол в формуле min(), снимая защиту, которую max(ttl_days, floor) давал раньше), ИЛИ cap_mult < 1 (тот же класс дыры, но со стороны потолка, а не пола: cap_mult приходит из jsonb default_params расписания -- ЕДИНСТВЕННЫЙ запланированный способ его задать, т.е. именно там опечатка 0 / 0.5 вместо 6 доходит до прода. cap_mult=0 даёт capped=0 -> effective_ttl_days=0 -> UPDATE снимает практически весь активный пул источника; cap_mult<1 (например 0.5) опускает потолок НИЖЕ заданного оператором ttl_days -- прямое нарушение инварианта «потолок не может понизить TTL ниже настроенного», который проверяет test_cap_never_lowers_ttl_below_configured_value), ИЛИ ttl_days/cap_mult -- bool (найдено ревью круга 3, 2026-08-15: `cap_mult < 1` пропускает `True` -- `bool` наследует `int`, `True < 1` ложно, а `ttl_days * True` == `ttl_days`, то есть потолок = сам ttl_days и пол молча отключается, никакого ValueError. jsonb `true`/`false` вместо числа -- ровно та опечатка в расписании, ради которой оба guard'а вообще написаны, поэтому bool отклоняется явной type-проверкой ДО числового сравнения для обоих параметров), ЛИБО если заданы одновременно null_segment_only=True и segments (взаимоисключающие срезы -- IS NULL и ANY(:segments) не композируются), ЛИБО (PR-B) min_floor_pairs < 0 / floor_drop_ratio < 1 / max_deactivated <= 0, ИЛИ любой из этих трёх -- bool (тот же класс jsonb-опечатки true/false вместо числа, что и у ttl_days/cap_mult выше -- default_params расписания это единственный запланированный способ их переопределить). """ counters: dict[str, int] = {"deactivated": 0} try: # bool -- подкласс int в Python, поэтому `True < 1` (False) и `False <= 0` # (True) НЕ ловят опечатку `"ttl_days": true` / `"cap_mult": true` в jsonb: # `ttl_days * True` == `ttl_days`, `cap_mult=True` даёт потолок == ttl_days и # молча отключает пол (см. Raises выше). Проверка типа -- ДО числового # сравнения, иначе bool проскакивает мимо него необнаруженным. if isinstance(ttl_days, bool): raise ValueError(f"ttl_days must be a number, not bool: {ttl_days!r}") if ttl_days <= 0: raise ValueError(f"ttl_days must be positive, got {ttl_days!r}") # Тот же класс дыры, что и ttl_days<=0 выше, только со стороны потолка: # cap_mult < 1 может опустить потолок (ttl_days * cap_mult) НИЖЕ заданного # ttl_days, а cap_mult <= 0 -- сделать капнутый потолок <= 0 и победить пол # в min() молча (ровно та дыра, ради которой заведён guard выше). Единственный # запланированный способ задать cap_mult -- вписать его руками в jsonb # default_params расписания (см. миграцию для avito), т.е. именно там опечатка # 0 / 0.5 вместо 6 -- реальный риск, а не гипотетика. if isinstance(cap_mult, bool): raise ValueError(f"cap_mult must be a number, not bool: {cap_mult!r}") if cap_mult < 1: raise ValueError(f"cap_mult must be >= 1, got {cap_mult!r}") # PR-B (#2659 продолжение): те же jsonb-bool-опечатки, тот же класс дыры, # что у ttl_days/cap_mult выше -- min_floor_pairs/floor_drop_ratio/ # max_deactivated тоже приходят из default_params расписания. if isinstance(min_floor_pairs, bool): raise ValueError(f"min_floor_pairs must be a number, not bool: {min_floor_pairs!r}") if min_floor_pairs < 0: raise ValueError(f"min_floor_pairs must be >= 0, got {min_floor_pairs!r}") if isinstance(floor_drop_ratio, bool): raise ValueError( f"floor_drop_ratio must be a number, not bool: {floor_drop_ratio!r}" ) if floor_drop_ratio < 1: raise ValueError(f"floor_drop_ratio must be >= 1, got {floor_drop_ratio!r}") if isinstance(max_deactivated, bool): raise ValueError(f"max_deactivated must be a number, not bool: {max_deactivated!r}") if max_deactivated <= 0: raise ValueError(f"max_deactivated must be positive, got {max_deactivated!r}") # Whitelist-проверка ДО построения/выполнения SQL: только после неё имя колонки # интерполируется f-string'ом. Значения по-прежнему идут через param-binding. # Внутри try -> невалидная колонка финализирует run как failed (mark_failed), # а не оставляет его «running» (единый контракт с прочими сбоями). if staleness_column not in _ALLOWED_STALENESS_COLUMNS: raise ValueError( f"invalid staleness_column={staleness_column!r}; " f"allowed: {sorted(_ALLOWED_STALENESS_COLUMNS)}" ) # null_segment_only + segments одновременно -- неоднозначный запрос: # IS NULL и ANY(:segments) -- разные, непересекающиеся предикаты, а не # композиция. Явный ValueError лучше молчаливого выбора одного из двух. if null_segment_only and segments is not None: raise ValueError("null_segment_only=True несовместимо с заданным segments") # Гейт по здоровью сбора (#2659) — ДО любого UPDATE. Деактивация необратима # на практике (вернуть «живость» может только повторный сбор), поэтому # проверяем ПЕРЕД записью, а не откатываем после. if min_confirmations > 0: health_params: dict[str, Any] = { "listing_source": listing_source, "health_window_days": health_window_days, } if segments is not None: health_params["segments"] = segments confirmations = ( db.execute( _build_confirmations_sql( staleness_column, with_segments=segments is not None, null_segment_only=null_segment_only, ), health_params, ).scalar() or 0 ) counters["confirmations"] = int(confirmations) if confirmations < min_confirmations: counters["skipped_unhealthy"] = 1 # Ничего не писали (был только SELECT) — rollback закрывает транзакцию # чисто, чтобы mark_done стартовал со своей. db.rollback() runs_mod.mark_done(db, run_id, counters) logger.warning( "deactivate_stale source=%s run_id=%d SKIPPED: сбор нездоров — " "подтверждений за %d сут %d < порога %d " "(segments=%r, null_segment_only=%s, staleness_column=%s); " "ни одна строка не тронута", listing_source, run_id, health_window_days, confirmations, min_confirmations, segments, null_segment_only, staleness_column, ) return counters # Пол TTL по измеренному циклу переобхода (#2659) — тоже ДО UPDATE и по тому же # срезу. Поднимает порог (max), но не выше потолка cap_mult * ttl_days (min) — # см. комментарий у CAP_MULT про петлю с положительной обратной связью и про # то, почему cap_mult -- параметр, а не голая константа. effective_ttl_days = ttl_days if revisit_floor_quantile > 0: floor_params: dict[str, Any] = { "listing_source": listing_source, "health_window_days": health_window_days, "revisit_quantile": revisit_floor_quantile, } if segments is not None: floor_params["segments"] = segments floor_days = db.execute( _build_revisit_floor_sql( staleness_column, with_segments=segments is not None, null_segment_only=null_segment_only, ), floor_params, ).scalar() # Размер выборки, из которой percentile_disc выше берёт квантиль (та же # общая _revisit_floor_from_where_sql). Появилось в PR-A как чистое # наблюдение; с PR-B (#2659 продолжение) гейтит деградацию -- см. блок # ниже. n_pairs обязан попасть в counters ДО того, как схлопнувшуюся # выборку станет видно только по повторению прод-инцидента. Считается # ВСЕГДА при включённом полу, в том числе когда floor_days получится # NULL -- тогда n_pairs=0 и объясняет, почему пол не посчитался. floor_n_pairs = db.execute( _build_revisit_floor_pairs_count_sql( staleness_column, with_segments=segments is not None, null_segment_only=null_segment_only, ), floor_params, ).scalar() counters["floor_n_pairs"] = int(floor_n_pairs or 0) # ── Гейт деградации пола (PR-B, #2659 продолжение) ──────────────── # ДО того, как floor_days (если есть) поднимет TTL -- вырожденная # выборка не должна ни давать "уверенный" пол, ни (тем более) # пропускать UPDATE вовсе. См. комментарий у DEFAULT_MIN_FLOOR_PAIRS / # DEFAULT_FLOOR_DROP_RATIO про то, почему сравнение идёт "прогон с # собой", а не с заново калиброванным порогом. prev_floor_n_pairs = db.execute( _PREVIOUS_FLOOR_N_PAIRS_SQL, {"run_id": run_id} ).scalar() below_min_floor_pairs = counters["floor_n_pairs"] < min_floor_pairs dropped_vs_previous = False if prev_floor_n_pairs is not None: counters["floor_n_pairs_previous"] = int(prev_floor_n_pairs) dropped_vs_previous = ( prev_floor_n_pairs > 0 and counters["floor_n_pairs"] * floor_drop_ratio < prev_floor_n_pairs ) if below_min_floor_pairs or dropped_vs_previous: counters["skipped_floor_degraded"] = 1 # Ничего не писали (только SELECT'ы) -- rollback закрывает # транзакцию чисто, чтобы mark_done стартовал со своей (тот же # приём, что у skipped_unhealthy выше). db.rollback() runs_mod.mark_done(db, run_id, counters) logger.warning( "deactivate_stale source=%s run_id=%d SKIPPED: пол деградировал -- " "floor_n_pairs=%d (порог %d), предыдущий успешный прогон=%s " "(порог падения ×%.1f) -- ни одна строка не тронута", listing_source, run_id, counters["floor_n_pairs"], min_floor_pairs, prev_floor_n_pairs, floor_drop_ratio, ) return counters # NULL = истории снимков за окно нет вовсе (свежая БД, дыра в снимках). # Тогда пола нет и TTL остаётся как задан: выдумывать пол не из чего. if floor_days is not None: counters["revisit_floor_days"] = ceil(float(floor_days)) # Пол поднимает TTL (max), потолок cap_mult его не пускает выше # ttl_days * cap_mult (min) — без этого пол растёт без ограничения # (см. комментарий у CAP_MULT). capped_ttl_days может быть float, # если cap_mult переопределён нецелым значением из default_params — # effective_ttl_days приводим к int (UPDATE ждёт целые сутки). raw_effective_ttl_days = max(ttl_days, counters["revisit_floor_days"]) capped_ttl_days = ttl_days * cap_mult effective_ttl_days = int(min(raw_effective_ttl_days, capped_ttl_days)) counters["ttl_days_effective"] = effective_ttl_days if raw_effective_ttl_days > capped_ttl_days: # Пол упёрся в потолок -- оба числа в counters (не только в логе), # чтобы это было видно в витрине прогонов, а не только в логах. # 1, а не True -- counters типизирован dict[str, int] (тот же # идиом, что skipped_unhealthy выше). counters["ttl_floor_capped"] = 1 counters["ttl_days_floor_raw"] = raw_effective_ttl_days logger.warning( "deactivate_stale source=%s run_id=%d TTL пол упёрся в потолок " "cap_mult=%s: пол поднял бы TTL до %d сут, потолок ограничивает " "заданные %d сут значением %d (квантиль %.3f, segments=%r) — " "растущий без ограничения пол это петля с положительной обратной " "связью, см. комментарий у CAP_MULT", listing_source, run_id, cap_mult, raw_effective_ttl_days, ttl_days, effective_ttl_days, revisit_floor_quantile, segments, ) elif effective_ttl_days > ttl_days: logger.warning( "deactivate_stale source=%s run_id=%d TTL поднят с %d до %d сут: " "свип за %d сут доказал живой строку, молчавшую %d сут " "(квантиль %.3f, segments=%r) — при ttl_days=%d снятие означало бы " "«мы не дошли», а не «объявление снято»", listing_source, run_id, ttl_days, effective_ttl_days, health_window_days, counters["revisit_floor_days"], revisit_floor_quantile, segments, ttl_days, ) # ── Потолок объёма снятия -- аварийный, наблюдательный (PR-B) ────── # Preflight count(*) ПО ТОМУ ЖЕ предикату, что исполнит UPDATE ниже # (тот же эффективный TTL, тот же срез сегментов) -- см. комментарий у # DEFAULT_MAX_DEACTIVATED про то, почему нет порога по ДОЛЕ пула # (метрика только пишется, не гейтит) и откуда взят абсолютный # аварийный порог. Живёт внутри revisit_floor_quantile > 0 -- на проде # это верно для всех активных расписаний деактивации, кроме # null_segment_only-джоб (миграция 266, пул на 2-3 порядка меньше # порога -- см. докстринг deactivate_stale_listings). preflight_params: dict[str, Any] = { "listing_source": listing_source, "ttl_days": effective_ttl_days, } if null_segment_only: candidates_sql = _build_null_segment_candidates_count_sql(staleness_column) elif segments is not None: candidates_sql = _build_segments_candidates_count_sql(staleness_column) preflight_params["segments"] = segments else: candidates_sql = _build_all_segments_candidates_count_sql(staleness_column) candidates_result = db.execute(candidates_sql, preflight_params).scalar() deactivation_candidates = int(candidates_result or 0) active_pool_result = db.execute( _ACTIVE_POOL_SQL, {"listing_source": listing_source} ).scalar() active_pool = int(active_pool_result or 0) counters["deactivation_candidates"] = deactivation_candidates counters["active_pool"] = active_pool counters["deactivated_pct"] = ( round(deactivation_candidates / active_pool * 100) if active_pool else 0 ) if deactivation_candidates > max_deactivated: counters["skipped_cap_exceeded"] = 1 db.rollback() runs_mod.mark_done(db, run_id, counters) logger.error( "deactivate_stale source=%s run_id=%d SKIPPED: кандидатов на снятие " "%d превышает аварийный потолок %d (active_pool=%d, %d%%) -- " "ни одна строка не тронута", listing_source, run_id, deactivation_candidates, max_deactivated, active_pool, counters["deactivated_pct"], ) return counters # null_segment_only -> IS NULL, отдельный явный предикат (ANY(:segments) # никогда не матчит NULL). segments is None -> все сегменты (поведение avito). # segments=[...] -> только перечисленные сегменты. Используем `is not None` # (НЕ truthy): пустой список [] означает "ни один сегмент" (= ANY(ARRAY[]) # ничего не матчит, деактивирует 0), а НЕ "все сегменты" — иначе случайный [] # стёр бы весь источник. if null_segment_only: params: dict[str, Any] = { "listing_source": listing_source, "ttl_days": effective_ttl_days, "run_id": run_id, } result = db.execute(_build_null_segment_sql(staleness_column), params) elif segments is not None: params = { "listing_source": listing_source, "ttl_days": effective_ttl_days, "segments": segments, "run_id": run_id, } result = db.execute(_build_segments_sql(staleness_column), params) else: params = { "listing_source": listing_source, "ttl_days": effective_ttl_days, "run_id": run_id, } result = db.execute(_build_all_segments_sql(staleness_column), params) counters["deactivated"] = result.rowcount or 0 db.commit() runs_mod.mark_done(db, run_id, counters) logger.info( "deactivate_stale source=%s run_id=%d done: deactivated=%d " "(ttl_days=%d эффективный, задан %d, segments=%r, null_segment_only=%s, " "staleness_column=%s)", listing_source, run_id, counters["deactivated"], effective_ttl_days, ttl_days, segments, null_segment_only, staleness_column, ) return counters except Exception as exc: logger.exception("deactivate_stale source=%s run_id=%d failed", listing_source, run_id) db.rollback() runs_mod.mark_failed(db, run_id, str(exc)[:1000], counters) raise def deactivate_stale_avito_listings(db: Session, run_id: int) -> dict[str, int]: """Обёртка для обратной совместимости -- deactivate_stale_avito (все сегменты, TTL из settings). Пометить is_active=false все avito-объявления с last_seen_at > TTL дней назад. Sync (вызывается scheduler-триггером в executor, как snapshot_listing_sources). Один UPDATE в транзакции. Финализирует scrape_runs (mark_done / mark_failed). Returns {"deactivated": N} -- количество обновлённых строк. """ return deactivate_stale_listings( db, run_id, listing_source="avito", ttl_days=settings.avito_stale_ttl_days, segments=None, )