"""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,
)