"""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 активных первичных строк) и NULL-сегмент не трогаем. - avito: все сегменты (segments=None), TTL=10 дней -- поведение без изменений. - Строки НЕ удаляются -- история нужна для бэктеста (#667). - #2674: деактивация в той же транзакции пишет снимок listings_snapshots со статусом 'stale' за текущую дату -- «мы N суток не видели». Жёсткое 'closed' (площадка ответила 404) пишет только avito_detail_backfill: смешивать факт с догадкой дорого. Задача синхронная (DB-only, никаких внешних HTTP-вызовов) -- запускается kit-scheduler'ом через product_handlers._job_deactivate_stale (wildcard-handler deactivate_stale_*), по образцу 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 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[]))" def _build_confirmations_sql(staleness_column: str, *, with_segments: bool) -> Any: """SELECT count(*) подтверждённых за окно строк — тот же срез, что и у UPDATE. staleness_column уже прошёл whitelist-проверку в deactivate_stale_listings. Значения (:listing_source, :health_window_days, :segments) — param-binding, psycopg v3 safe (CAST(... AS ...), никаких :param::type). """ segment_filter = _CONFIRMATIONS_SEGMENT_FILTER if with_segments else "" 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} """ ) # Дефолтные (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 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, ) -> dict[str, int]: """Пометить is_active=false объявления, чья свежесть старше ttl_days дней. Параметры: listing_source: значение поля source в таблице listings ('avito', 'yandex', 'cian', 'domklik'). ttl_days: количество дней TTL; объявления старше этого порога деактивируются. segments: если задан -- деактивировать только объявления с указанными listing_segment значениями. None -> все сегменты (поведение avito по умолчанию). 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. 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} и НИ ОДНА строка не тронута. Raises: ValueError: если staleness_column не входит в whitelist (проверка ДО SQL, никакой интерполяции пользовательского ввода в запрос). """ counters: dict[str, int] = {"deactivated": 0} try: # 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)}" ) # Гейт по здоровью сбора (#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), 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, staleness_column=%s); ни одна строка не тронута", listing_source, run_id, health_window_days, confirmations, min_confirmations, segments, staleness_column, ) return counters # segments is None -> все сегменты (поведение avito). segments=[...] -> только # перечисленные сегменты. Используем `is not None` (НЕ truthy): пустой список [] # означает "ни один сегмент" (= ANY(ARRAY[]) ничего не матчит, деактивирует 0), # а НЕ "все сегменты" — иначе случайный [] стёр бы весь источник. if segments is not None: params: dict[str, Any] = { "listing_source": listing_source, "ttl_days": ttl_days, "segments": segments, "run_id": run_id, } result = db.execute(_build_segments_sql(staleness_column), params) else: params = { "listing_source": listing_source, "ttl_days": ttl_days, "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, segments=%r, staleness_column=%s)", listing_source, run_id, counters["deactivated"], ttl_days, segments, 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, )