From 22ccc1050cb7e42f64ae368188391a96a333e966 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Mon, 24 Aug 2026 00:24:56 +0300 Subject: [PATCH] =?UTF-8?q?feat(tradein/deactivate-stale):=20=D1=81=D1=82?= =?UTF-8?q?=D1=80=D0=B0=D1=85=D0=BE=D0=B2=D0=BE=D1=87=D0=BD=D1=8B=D0=B5=20?= =?UTF-8?q?=D1=80=D0=B5=D0=BB=D1=8C=D1=81=D1=8B=20=D0=BE=D0=B1=D1=8A=D1=91?= =?UTF-8?q?=D0=BC=D0=B0=20=D1=81=D0=BD=D1=8F=D1=82=D0=B8=D1=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Продолжение #2659 после PR-A. У джобы деактивации не было НИ ОДНОГО ограничителя ОБЪЁМА: min_confirmations может пропустить прогон целиком, пол переобхода поднимает TTL, cap_mult ограничивает сам пол — но ни один не смотрит на РАЗМЕР ВЫБОРКИ, из которой посчитан квантиль, и ни один не смотрит на то, сколько строк UPDATE снимет. UPDATE идёт одним statement'ом без LIMIT: при обвале любого гейта выше снимается сколько снимется. Добавлено два рельса, оба внутри revisit_floor_quantile > 0: 1. Гейт деградации пола. floor_n_pairs (счётчик из PR-A) ниже min_floor_pairs=30 ЛИБО упал более чем впятеро относительно предыдущего успешного прогона ТОГО ЖЕ РАСПИСАНИЯ -> skipped_floor_degraded, ни одна строка не тронута. Сравнение идёт «прогон сам с собой» по scrape_runs.source, а не по listing_source: у cian и cian_null_segment listing_source один, но это две серии с несопоставимым масштабом выборки — сравнивать их было бы категориальной ошибкой. 2. Аварийный потолок max_deactivated=15000. Preflight count(*) по ТОМУ ЖЕ предикату, что исполнит UPDATE — abort ДО записи. Не рабочий порог, а предохранитель: исторический максимум легитимного снятия 9 300 (avito 06.06), запас 1.6x. Порога по ДОЛЕ пула сознательно НЕТ: пул, с которым джоба реально работает, никогда не записывался, а восстановленный постфактум ряд даёт 40 %, 97 %, 317 %, 1652 % — калибровать не на чем. Поэтому deactivation_candidates, active_pool и deactivated_pct пока только ПИШУТСЯ. Валидация параметров — тот же класс, что у ttl_days/cap_mult: bool отбивается ДО числового сравнения (jsonb-опечатка true/false в default_params расписания — единственный запланированный способ их задать). uv.lock: две строки requires-dist подтянуты к pyproject (curl-cffi >=0.7.0 -> >=0.15.0). Это устранение устаревшей записи, а не смена зависимости — разрешённые версии не изменились; лок расходился с pyproject на main и регенерировался при любом вызове uv. Проверено: 249 passed, 1 skipped (-k "deactivate or stale"). Refs #2659 --- .../backend/app/services/product_handlers.py | 14 + .../app/tasks/deactivate_stale_avito.py | 345 +++++++++++++++++- .../test_deactivate_stale_deactivation_cap.py | 311 ++++++++++++++++ ...test_deactivate_stale_floor_degradation.py | 317 ++++++++++++++++ tradein-mvp/uv.lock | 4 +- 5 files changed, 976 insertions(+), 15 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_deactivate_stale_deactivation_cap.py create mode 100644 tradein-mvp/backend/tests/test_deactivate_stale_floor_degradation.py diff --git a/tradein-mvp/backend/app/services/product_handlers.py b/tradein-mvp/backend/app/services/product_handlers.py index f22dd512..4076fa67 100644 --- a/tradein-mvp/backend/app/services/product_handlers.py +++ b/tradein-mvp/backend/app/services/product_handlers.py @@ -218,7 +218,10 @@ async def _job_deactivate_stale( from app.core.config import settings as _settings from app.tasks.deactivate_stale_avito import ( CAP_MULT, + DEFAULT_FLOOR_DROP_RATIO, + DEFAULT_MAX_DEACTIVATED, DEFAULT_MIN_CONFIRMATIONS, + DEFAULT_MIN_FLOOR_PAIRS, DEFAULT_REVISIT_FLOOR_QUANTILE, deactivate_stale_listings, ) @@ -247,6 +250,14 @@ async def _job_deactivate_stale( # относительно своего ttl_days переопределяет его через default_params (ключ # "cap_mult"), не трогая дефолт для остальных источников. cap_mult: float = params.get("cap_mult", CAP_MULT) + # Гейт деградации пола (PR-B, #2659 продолжение) -- включён по умолчанию, тот + # же принцип, что у min_confirmations/revisit_floor_quantile выше: незасеянное + # расписание получает страховку, а не «деактивируй вслепую». + min_floor_pairs: int = params.get("min_floor_pairs", DEFAULT_MIN_FLOOR_PAIRS) + floor_drop_ratio: float = params.get("floor_drop_ratio", DEFAULT_FLOOR_DROP_RATIO) + # Аварийный (не рабочий) потолок объёма снятия за один прогон -- см. + # DEFAULT_MAX_DEACTIVATED в deactivate_stale_avito.py. + max_deactivated: int = params.get("max_deactivated", DEFAULT_MAX_DEACTIVATED) loop = asyncio.get_event_loop() await loop.run_in_executor( @@ -262,6 +273,9 @@ async def _job_deactivate_stale( revisit_floor_quantile=revisit_floor_quantile, null_segment_only=null_segment_only, cap_mult=cap_mult, + min_floor_pairs=min_floor_pairs, + floor_drop_ratio=floor_drop_ratio, + max_deactivated=max_deactivated, ), ) diff --git a/tradein-mvp/backend/app/tasks/deactivate_stale_avito.py b/tradein-mvp/backend/app/tasks/deactivate_stale_avito.py index 5113649e..38834eb8 100644 --- a/tradein-mvp/backend/app/tasks/deactivate_stale_avito.py +++ b/tradein-mvp/backend/app/tasks/deactivate_stale_avito.py @@ -276,6 +276,61 @@ _REVISIT_FLOOR_NULL_SEGMENT_FILTER = "\n AND l.listing_segment IS NULL" 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. @@ -513,6 +568,122 @@ _DEACTIVATE_SQL_SEGMENTS = _build_segments_sql("last_seen_at") _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, @@ -526,6 +697,9 @@ def deactivate_stale_listings( 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 дней. @@ -568,20 +742,52 @@ def deactivate_stale_listings( источники. Итоговый потолок = 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, + Если гейт здоровья не пропустил прогон: {"deactivated": 0, "confirmations": N, "skipped_unhealthy": 1} и НИ ОДНА строка не тронута. Если пол переобхода включён (revisit_floor_quantile > 0): дополнительно {"floor_n_pairs": N} -- размер - выборки, из которой percentile_disc посчитал квантиль (наблюдательность, PR-A - #2659 продолжение; сейчас ничего не гейтит). Если пол при этом реально поднял - TTL: дополнительно {"revisit_floor_days": N, "ttl_days_effective": N}. Если пол - упёрся в потолок cap_mult: дополнительно {"ttl_floor_capped": 1, - "ttl_days_floor_raw": N} -- N это то, во что пол поднял бы TTL БЕЗ потолка. + выборки, из которой 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 @@ -605,7 +811,11 @@ def deactivate_stale_listings( которой оба guard'а вообще написаны, поэтому bool отклоняется явной type-проверкой ДО числового сравнения для обоих параметров), ЛИБО если заданы одновременно null_segment_only=True и segments (взаимоисключающие - срезы -- IS NULL и ANY(: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: @@ -631,6 +841,26 @@ def deactivate_stale_listings( 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), @@ -711,12 +941,13 @@ def deactivate_stale_listings( ), floor_params, ).scalar() - # Наблюдательность (PR-A, #2659 продолжение) -- та же выборка, что и - # percentile_disc выше (общая _revisit_floor_from_where_sql). Пока ничего - # не гейтит (это PR-B), но n_pairs обязан попасть в counters ДО того, как - # схлопнувшуюся выборку станет видно только по повторению прод-инцидента. - # Считается ВСЕГДА при включённом полу, в том числе когда floor_days - # получится NULL -- тогда n_pairs=0 и объясняет, почему пол не посчитался. + # Размер выборки, из которой 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, @@ -726,6 +957,44 @@ def deactivate_stale_listings( 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: @@ -779,6 +1048,56 @@ def deactivate_stale_listings( 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` diff --git a/tradein-mvp/backend/tests/test_deactivate_stale_deactivation_cap.py b/tradein-mvp/backend/tests/test_deactivate_stale_deactivation_cap.py new file mode 100644 index 00000000..9fa7baee --- /dev/null +++ b/tradein-mvp/backend/tests/test_deactivate_stale_deactivation_cap.py @@ -0,0 +1,311 @@ +"""Потолок объёма снятия за один прогон -- аварийный, наблюдательный (PR-B, #2659 +продолжение). + +У джобы деактивации не было НИ ОДНОГО ограничителя ОБЪЁМА: UPDATE идёт одним +statement'ом без LIMIT, counters["deactivated"] = result.rowcount. Долевой порог +("не больше N% активного пула") откалибровать по историческим данным нельзя -- +исторические снятия дают доли 40%/97%/317%/1652% от восстановленного пула, числа +не образуют осмысленного ряда (гранулярность listings vs listing_sources разная). +Поэтому: preflight count(*) ПО ТОМУ ЖЕ предикату, что и UPDATE, пишет +deactivation_candidates/active_pool/deactivated_pct В КАЖДЫЙ прогон (наблюдение, +задел под будущую калибровку долевого порога), а блокирует ТОЛЬКО абсолютный +аварийный порог max_deactivated (дефолт 15000 -- запас 1.6x над историческим +максимумом легитимного снятия 9300, avito 2026-06-06). +""" + +from __future__ import annotations + +import os +from typing import Any + +import pytest + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +from app.tasks import deactivate_stale_avito as task_mod + +# ── Фейковая сессия (тот же контракт, что test_deactivate_stale_floor_degradation.py) ── + + +class _FakeResult: + def __init__(self, rowcount: int = 0, scalar_value: Any = None) -> None: + self.rowcount = rowcount + self._scalar = scalar_value + + def scalar(self) -> Any: + return self._scalar + + +class _FakeDB: + def __init__( + self, + *, + floor_days: float | None = 50.0, + floor_n_pairs: int = 5000, + prev_floor_n_pairs: int | None = 5000, + confirmations: int = 10_000, + deactivation_candidates: int | None = None, + active_pool: int = 100_000, + rowcount: int = 137, + ) -> None: + self._floor = floor_days + self._floor_n_pairs = floor_n_pairs + self._prev_floor_n_pairs = prev_floor_n_pairs + self._confirmations = confirmations + self._deactivation_candidates = ( + rowcount if deactivation_candidates is None else deactivation_candidates + ) + self._active_pool = active_pool + self._rowcount = rowcount + self.executed: list[tuple[str, dict[str, Any] | None]] = [] + self.committed = False + self.rolled_back = False + + def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult: + sql = str(stmt.text) + self.executed.append((sql, params)) + if "percentile_disc" in sql: + return _FakeResult(scalar_value=self._floor) + if "FROM scrape_runs prev" in sql: + return _FakeResult(scalar_value=self._prev_floor_n_pairs) + if "JOIN LATERAL" in sql: + return _FakeResult(scalar_value=self._floor_n_pairs) + if "health_window_days" in sql: + return _FakeResult(scalar_value=self._confirmations) + if "SELECT count(*)" in sql and "ttl_days" in sql: + return _FakeResult(scalar_value=self._deactivation_candidates) + if "SELECT count(*)" in sql: + return _FakeResult(scalar_value=self._active_pool) + return _FakeResult(rowcount=self._rowcount) + + def commit(self) -> None: + self.committed = True + + def rollback(self) -> None: + self.rolled_back = True + + @property + def update_query(self) -> tuple[str, dict[str, Any] | None]: + return next((e for e in self.executed if "UPDATE listings" in e[0]), ("", None)) + + @property + def executed_kinds(self) -> list[str]: + kinds = [] + for sql, _ in self.executed: + if "percentile_disc" in sql: + kinds.append("floor") + elif "FROM scrape_runs prev" in sql: + kinds.append("prev_floor_n_pairs") + elif "JOIN LATERAL" in sql: + kinds.append("floor_n_pairs") + elif "health_window_days" in sql: + kinds.append("confirmations") + elif "SELECT count(*)" in sql and "ttl_days" in sql: + kinds.append("candidates") + elif "SELECT count(*)" in sql: + kinds.append("active_pool") + elif "UPDATE listings" in sql: + kinds.append("update") + else: + kinds.append("unknown") + return kinds + + +def _run(db: _FakeDB, monkeypatch: pytest.MonkeyPatch, **kwargs: Any) -> dict[str, int]: + monkeypatch.setattr(task_mod.runs_mod, "mark_done", lambda *a, **k: None) + monkeypatch.setattr(task_mod.runs_mod, "mark_failed", lambda *a, **k: None) + kwargs.setdefault("revisit_floor_quantile", task_mod.DEFAULT_REVISIT_FLOOR_QUANTILE) + return task_mod.deactivate_stale_listings( + db, # type: ignore[arg-type] + 1, + listing_source=kwargs.pop("listing_source", "avito"), + ttl_days=kwargs.pop("ttl_days", 10), + **kwargs, + ) + + +# ── Срабатывание потолка ───────────────────────────────────────────────────────── + + +def test_cap_trips_when_candidates_exceed_max_deactivated(monkeypatch: pytest.MonkeyPatch) -> None: + db = _FakeDB(deactivation_candidates=20_000, active_pool=100_000) + out = _run(db, monkeypatch, max_deactivated=15000) + assert out["skipped_cap_exceeded"] == 1 + assert out["deactivation_candidates"] == 20000 + assert out["active_pool"] == 100000 + assert out["deactivated"] == 0 + assert db.update_query[0] == "" + assert db.committed is False + assert db.rolled_back is True + + +def test_cap_does_not_trip_below_threshold(monkeypatch: pytest.MonkeyPatch) -> None: + db = _FakeDB(deactivation_candidates=9300, active_pool=100_000, rowcount=9300) + out = _run(db, monkeypatch, max_deactivated=15000) + assert "skipped_cap_exceeded" not in out + assert out["deactivation_candidates"] == 9300 + assert out["deactivated"] == 9300 + assert db.update_query[0] != "" + assert db.committed is True + + +def test_cap_trips_exactly_above_threshold_not_at_it(monkeypatch: pytest.MonkeyPatch) -> None: + """== max_deactivated не триггерит, только > (строго больше).""" + db = _FakeDB(deactivation_candidates=15000, active_pool=100_000, rowcount=15000) + out = _run(db, monkeypatch, max_deactivated=15000) + assert "skipped_cap_exceeded" not in out + db2 = _FakeDB(deactivation_candidates=15001, active_pool=100_000) + out2 = _run(db2, monkeypatch, max_deactivated=15000) + assert out2["skipped_cap_exceeded"] == 1 + + +def test_cap_check_runs_before_update(monkeypatch: pytest.MonkeyPatch) -> None: + db = _FakeDB(deactivation_candidates=20_000) + _run(db, monkeypatch, max_deactivated=15000) + assert "update" not in db.executed_kinds + + +def test_blocked_run_is_finalised_as_done(monkeypatch: pytest.MonkeyPatch) -> None: + marked: dict[str, Any] = {} + monkeypatch.setattr( + task_mod.runs_mod, + "mark_done", + lambda _db, run_id, counters: marked.update(run_id=run_id, counters=dict(counters)), + ) + monkeypatch.setattr(task_mod.runs_mod, "mark_failed", lambda *a, **k: None) + db = _FakeDB(deactivation_candidates=20_000) + task_mod.deactivate_stale_listings( + db, # type: ignore[arg-type] + 55, + listing_source="avito", + ttl_days=10, + revisit_floor_quantile=task_mod.DEFAULT_REVISIT_FLOOR_QUANTILE, + max_deactivated=15000, + ) + assert marked["run_id"] == 55 + assert marked["counters"]["skipped_cap_exceeded"] == 1 + assert marked["counters"]["deactivated"] == 0 + + +# ── Наблюдательные counters (пишутся ВСЕГДА, не только при abort) ──────────────── + + +def test_deactivated_pct_is_computed(monkeypatch: pytest.MonkeyPatch) -> None: + db = _FakeDB(deactivation_candidates=5000, active_pool=10_000, rowcount=5000) + out = _run(db, monkeypatch) + assert out["deactivated_pct"] == 50 + + +def test_deactivated_pct_zero_when_active_pool_empty(monkeypatch: pytest.MonkeyPatch) -> None: + """active_pool=0 -> деление защищено, не ZeroDivisionError.""" + db = _FakeDB(deactivation_candidates=0, active_pool=0, rowcount=0) + out = _run(db, monkeypatch) + assert out["deactivated_pct"] == 0 + assert out["active_pool"] == 0 + + +def test_candidates_and_pool_written_on_every_healthy_run(monkeypatch: pytest.MonkeyPatch) -> None: + db = _FakeDB(deactivation_candidates=137, active_pool=5000, rowcount=137) + out = _run(db, monkeypatch) + assert out["deactivation_candidates"] == 137 + assert out["active_pool"] == 5000 + assert db.update_query[0] != "" + + +# ── Предикат preflight COUNT совпадает с предикатом UPDATE ─────────────────────── + + +def test_candidates_predicate_matches_update_predicate_all_segments() -> None: + update_sql = str(task_mod._build_all_segments_sql("last_seen_at").text) + count_sql = str(task_mod._build_all_segments_candidates_count_sql("last_seen_at").text) + for fragment in ( + "source = :listing_source", + "is_active = true", + "last_seen_at < NOW() - CAST(:ttl_days || ' days' AS interval)", + ): + assert fragment in update_sql, f"{fragment!r} missing from UPDATE sql" + assert fragment in count_sql, f"{fragment!r} missing from candidates-count sql" + assert "listing_segment" not in update_sql + assert "listing_segment" not in count_sql + + +def test_candidates_predicate_matches_update_predicate_segments() -> None: + update_sql = str(task_mod._build_segments_sql("last_seen_at").text) + count_sql = str(task_mod._build_segments_candidates_count_sql("last_seen_at").text) + for fragment in ( + "source = :listing_source", + "is_active = true", + "last_seen_at < NOW() - CAST(:ttl_days || ' days' AS interval)", + "listing_segment = ANY(CAST(:segments AS text[]))", + ): + assert fragment in update_sql, f"{fragment!r} missing from UPDATE sql" + assert fragment in count_sql, f"{fragment!r} missing from candidates-count sql" + + +def test_candidates_predicate_matches_update_predicate_null_segment() -> None: + update_sql = str(task_mod._build_null_segment_sql("last_seen_at").text) + count_sql = str(task_mod._build_null_segment_candidates_count_sql("last_seen_at").text) + for fragment in ( + "source = :listing_source", + "is_active = true", + "last_seen_at < NOW() - CAST(:ttl_days || ' days' AS interval)", + "listing_segment IS NULL", + ): + assert fragment in update_sql, f"{fragment!r} missing from UPDATE sql" + assert fragment in count_sql, f"{fragment!r} missing from candidates-count sql" + + +def test_candidates_count_sql_uses_the_effective_ttl_days_param_name() -> None: + """Preflight использует :ttl_days -- caller обязан передать effective_ttl_days + под этим именем (см. deactivate_stale_listings), а не сырой ttl_days.""" + sql = str(task_mod._build_all_segments_candidates_count_sql("last_seen_at").text) + assert ":ttl_days" in sql + + +def test_active_pool_sql_has_no_segment_or_ttl_filter() -> None: + """active_pool -- весь активный пул source, без сегмента и без TTL (см. докстринг + deactivate_stale_listings, Returns): денормализатор пула, а не срез UPDATE.""" + sql = str(task_mod._ACTIVE_POOL_SQL.text) + assert "is_active = true" in sql + assert "listing_segment" not in sql + assert ":ttl_days" not in sql + + +# ── Параметры / контракт ────────────────────────────────────────────────────────── + + +def test_default_max_deactivated_is_15000() -> None: + assert task_mod.DEFAULT_MAX_DEACTIVATED == 15000 + + +def test_default_does_not_reject_any_known_legit_historical_run() -> None: + """Исторический максимум легитимного снятия (avito, 2026-06-06) и следующие по + убыванию -- ни один не должен упереться в дефолтный потолок.""" + known_legit = [9300, 6531, 6131, 4959, 3909] + assert max(known_legit) < task_mod.DEFAULT_MAX_DEACTIVATED + + +def test_max_deactivated_rejects_non_positive(monkeypatch: pytest.MonkeyPatch) -> None: + db = _FakeDB() + with pytest.raises(ValueError, match="max_deactivated"): + _run(db, monkeypatch, max_deactivated=0) + assert db.executed == [] + + +def test_max_deactivated_rejects_bool(monkeypatch: pytest.MonkeyPatch) -> None: + db = _FakeDB() + with pytest.raises(ValueError, match="max_deactivated"): + _run(db, monkeypatch, max_deactivated=True) + assert db.executed == [] + + +def test_handler_wires_max_deactivated_from_schedule() -> None: + """Читаем исходник файлом (как соседние test_handler_wires_* в этом сьюте): + product_handlers тянет scraper_kit, которого в юнит-окружении может не быть.""" + from pathlib import Path + + handlers = Path(__file__).resolve().parents[1] / "app" / "services" / "product_handlers.py" + src = handlers.read_text("utf-8") + job = src.split("async def _job_deactivate_stale")[1].split("\nasync def ")[0] + assert 'params.get("max_deactivated", DEFAULT_MAX_DEACTIVATED)' in job + assert "max_deactivated=max_deactivated" in job diff --git a/tradein-mvp/backend/tests/test_deactivate_stale_floor_degradation.py b/tradein-mvp/backend/tests/test_deactivate_stale_floor_degradation.py new file mode 100644 index 00000000..3708f4b4 --- /dev/null +++ b/tradein-mvp/backend/tests/test_deactivate_stale_floor_degradation.py @@ -0,0 +1,317 @@ +"""Гейт деградации пола переобхода (PR-B, #2659 продолжение). + +PR-A (fb38d657) перевёл поиск предшественника в поле переобхода на LATERAL и +добавил counters["floor_n_pairs"] -- ЧИСТО наблюдательный счётчик размера +выборки, из которой percentile_disc берёт квантиль. Этот файл проверяет PR-B: +тот же floor_n_pairs теперь ГЕЙТИТ прогон, если выборка вырождена -- либо ниже +абсолютного порога (min_floor_pairs), либо упала более чем в floor_drop_ratio +раз относительно предыдущего УСПЕШНОГО прогона того же расписания. +""" + +from __future__ import annotations + +import os +import re +from typing import Any + +import pytest + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +from app.tasks import deactivate_stale_avito as task_mod + +# ── Фейковая сессия ───────────────────────────────────────────────────────────── +# Различает ВСЕ SELECT-варианты, которые видит deactivate_stale_listings при +# revisit_floor_quantile > 0: percentile_disc (пол), previous floor_n_pairs +# (scrape_runs, PR-B), floor_n_pairs (LATERAL count, PR-A), confirmations (гейт +# здоровья), deactivation_candidates / active_pool (потолок объёма, PR-B). +# +# Порядок веток важен: percentile_disc и count(*) LATERAL-запрос ОБА содержат +# "JOIN LATERAL" -- percentile_disc проверяется первой веткой. LATERAL-запрос +# также содержит "health_window_days" (дважды) -- "JOIN LATERAL" проверяется +# раньше этой ветки. UPDATE содержит "ttl_days", но не "SELECT count(*)" -- +# ветка candidates требует ОБА маркера, поэтому UPDATE в неё не попадает. + + +class _FakeResult: + def __init__(self, rowcount: int = 0, scalar_value: Any = None) -> None: + self.rowcount = rowcount + self._scalar = scalar_value + + def scalar(self) -> Any: + return self._scalar + + +class _FakeDB: + def __init__( + self, + *, + floor_days: float | None = 50.0, + floor_n_pairs: int = 5000, + prev_floor_n_pairs: int | None = 5000, + confirmations: int = 10_000, + deactivation_candidates: int | None = None, + active_pool: int = 100_000, + rowcount: int = 137, + ) -> None: + self._floor = floor_days + self._floor_n_pairs = floor_n_pairs + self._prev_floor_n_pairs = prev_floor_n_pairs + self._confirmations = confirmations + self._deactivation_candidates = ( + rowcount if deactivation_candidates is None else deactivation_candidates + ) + self._active_pool = active_pool + self._rowcount = rowcount + self.executed: list[tuple[str, dict[str, Any] | None]] = [] + self.committed = False + self.rolled_back = False + + def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult: + sql = str(stmt.text) + self.executed.append((sql, params)) + if "percentile_disc" in sql: + return _FakeResult(scalar_value=self._floor) + if "FROM scrape_runs prev" in sql: + return _FakeResult(scalar_value=self._prev_floor_n_pairs) + if "JOIN LATERAL" in sql: + return _FakeResult(scalar_value=self._floor_n_pairs) + if "health_window_days" in sql: + return _FakeResult(scalar_value=self._confirmations) + if "SELECT count(*)" in sql and "ttl_days" in sql: + return _FakeResult(scalar_value=self._deactivation_candidates) + if "SELECT count(*)" in sql: + return _FakeResult(scalar_value=self._active_pool) + return _FakeResult(rowcount=self._rowcount) + + def commit(self) -> None: + self.committed = True + + def rollback(self) -> None: + self.rolled_back = True + + @property + def update_query(self) -> tuple[str, dict[str, Any] | None]: + return next((e for e in self.executed if "UPDATE listings" in e[0]), ("", None)) + + @property + def executed_kinds(self) -> list[str]: + kinds = [] + for sql, _ in self.executed: + if "percentile_disc" in sql: + kinds.append("floor") + elif "FROM scrape_runs prev" in sql: + kinds.append("prev_floor_n_pairs") + elif "JOIN LATERAL" in sql: + kinds.append("floor_n_pairs") + elif "health_window_days" in sql: + kinds.append("confirmations") + elif "SELECT count(*)" in sql and "ttl_days" in sql: + kinds.append("candidates") + elif "SELECT count(*)" in sql: + kinds.append("active_pool") + elif "UPDATE listings" in sql: + kinds.append("update") + else: + kinds.append("unknown") + return kinds + + +def _run(db: _FakeDB, monkeypatch: pytest.MonkeyPatch, **kwargs: Any) -> dict[str, int]: + monkeypatch.setattr(task_mod.runs_mod, "mark_done", lambda *a, **k: None) + monkeypatch.setattr(task_mod.runs_mod, "mark_failed", lambda *a, **k: None) + kwargs.setdefault("revisit_floor_quantile", task_mod.DEFAULT_REVISIT_FLOOR_QUANTILE) + return task_mod.deactivate_stale_listings( + db, # type: ignore[arg-type] + 1, + listing_source=kwargs.pop("listing_source", "avito"), + ttl_days=kwargs.pop("ttl_days", 10), + **kwargs, + ) + + +# ── Абсолютный порог ──────────────────────────────────────────────────────────── + + +def test_gate_trips_on_absolute_floor(monkeypatch: pytest.MonkeyPatch) -> None: + """floor_n_pairs ниже min_floor_pairs -> skipped_floor_degraded, ни одна строка не тронута.""" + db = _FakeDB(floor_n_pairs=10, prev_floor_n_pairs=None) + out = _run(db, monkeypatch, min_floor_pairs=30) + assert out["skipped_floor_degraded"] == 1 + assert out["floor_n_pairs"] == 10 + assert out["deactivated"] == 0 + assert db.update_query[0] == "" + assert db.committed is False + assert db.rolled_back is True + + +def test_absolute_floor_disabled_at_zero(monkeypatch: pytest.MonkeyPatch) -> None: + """min_floor_pairs=0 -> абсолютная проверка отключена (тот же идиом, что у гейта здоровья).""" + db = _FakeDB(floor_n_pairs=0, prev_floor_n_pairs=None) + out = _run(db, monkeypatch, min_floor_pairs=0) + assert "skipped_floor_degraded" not in out + assert db.update_query[0] != "" + + +# ── Относительный порог (падение в N раз) ─────────────────────────────────────── + + +def test_gate_trips_on_ratio_drop(monkeypatch: pytest.MonkeyPatch) -> None: + """floor_n_pairs упал в 10 раз против предыдущего успешного прогона (порог 5x).""" + db = _FakeDB(floor_n_pairs=100, prev_floor_n_pairs=1000) + out = _run(db, monkeypatch, min_floor_pairs=30, floor_drop_ratio=5.0) + assert out["skipped_floor_degraded"] == 1 + assert out["floor_n_pairs_previous"] == 1000 + assert out["deactivated"] == 0 + assert db.update_query[0] == "" + assert db.committed is False + assert db.rolled_back is True + + +def test_gate_does_not_trip_at_exactly_the_ratio(monkeypatch: pytest.MonkeyPatch) -> None: + """Падение РОВНО в floor_drop_ratio раз не триггерит -- задача требует «более чем».""" + db = _FakeDB(floor_n_pairs=200, prev_floor_n_pairs=1000) # ровно 5x + out = _run(db, monkeypatch, min_floor_pairs=30, floor_drop_ratio=5.0) + assert "skipped_floor_degraded" not in out + assert db.update_query[0] != "" + + +def test_gate_ignores_ratio_when_no_previous_run(monkeypatch: pytest.MonkeyPatch) -> None: + """Нет предыдущего успешного прогона (первый прогон после деплоя) -> относительная + часть молчит, работает только абсолютный порог.""" + db = _FakeDB(floor_n_pairs=50, prev_floor_n_pairs=None) + out = _run(db, monkeypatch, min_floor_pairs=30, floor_drop_ratio=5.0) + assert "skipped_floor_degraded" not in out + assert "floor_n_pairs_previous" not in out + assert db.update_query[0] != "" + + +def test_gate_ignores_zero_previous(monkeypatch: pytest.MonkeyPatch) -> None: + """Предыдущий прогон существует, но floor_n_pairs=0 в нём -- деление защищено guard'ом + prev_floor_n_pairs > 0, ratio-часть молчит вместо ZeroDivisionError.""" + db = _FakeDB(floor_n_pairs=0, prev_floor_n_pairs=0) + out = _run(db, monkeypatch, min_floor_pairs=0, floor_drop_ratio=5.0) + assert "skipped_floor_degraded" not in out + assert out["floor_n_pairs_previous"] == 0 + + +# ── Нормальный (не деградировавший) floor_n_pairs ─────────────────────────────── + + +def test_gate_passes_with_healthy_floor_n_pairs(monkeypatch: pytest.MonkeyPatch) -> None: + db = _FakeDB(floor_n_pairs=5000, prev_floor_n_pairs=5200) + out = _run(db, monkeypatch, min_floor_pairs=30, floor_drop_ratio=5.0) + assert "skipped_floor_degraded" not in out + assert out["floor_n_pairs"] == 5000 + assert out["floor_n_pairs_previous"] == 5200 + assert db.update_query[0] != "" + assert db.committed is True + + +def test_gate_disabled_when_revisit_floor_is_off(monkeypatch: pytest.MonkeyPatch) -> None: + """revisit_floor_quantile=0 -> весь блок (пол + оба гейта PR-B) не исполняется вовсе.""" + db = _FakeDB(floor_n_pairs=1, prev_floor_n_pairs=None) + out = _run(db, monkeypatch, revisit_floor_quantile=0, min_floor_pairs=30) + assert out == {"deactivated": 137} + assert "skipped_floor_degraded" not in out + assert db.executed_kinds == ["update"] + + +# ── Abort -- ни одна строка не тронута ─────────────────────────────────────────── + + +def test_gate_runs_before_update_and_before_the_volume_cap(monkeypatch: pytest.MonkeyPatch) -> None: + """Деградация останавливает прогон ДО preflight-потолка (PR-B п.2) и ДО UPDATE.""" + db = _FakeDB(floor_n_pairs=10, prev_floor_n_pairs=None) + _run(db, monkeypatch, min_floor_pairs=30) + kinds = db.executed_kinds + assert "update" not in kinds + assert "candidates" not in kinds + assert "active_pool" not in kinds + + +def test_blocked_run_is_finalised_as_done(monkeypatch: pytest.MonkeyPatch) -> None: + """Пропущенный прогон закрывается mark_done, а не висит 'running' до zombie-жатвы.""" + marked: dict[str, Any] = {} + monkeypatch.setattr( + task_mod.runs_mod, + "mark_done", + lambda _db, run_id, counters: marked.update(run_id=run_id, counters=dict(counters)), + ) + monkeypatch.setattr(task_mod.runs_mod, "mark_failed", lambda *a, **k: None) + db = _FakeDB(floor_n_pairs=10, prev_floor_n_pairs=None) + task_mod.deactivate_stale_listings( + db, # type: ignore[arg-type] + 99, + listing_source="avito", + ttl_days=10, + revisit_floor_quantile=task_mod.DEFAULT_REVISIT_FLOOR_QUANTILE, + min_floor_pairs=30, + ) + assert marked["run_id"] == 99 + assert marked["counters"]["skipped_floor_degraded"] == 1 + assert marked["counters"]["deactivated"] == 0 + + +# ── Параметры / контракт ───────────────────────────────────────────────────────── + + +def test_default_min_floor_pairs_is_a_safety_net_not_zero() -> None: + assert task_mod.DEFAULT_MIN_FLOOR_PAIRS > 0 + + +def test_default_floor_drop_ratio_is_above_one() -> None: + assert task_mod.DEFAULT_FLOOR_DROP_RATIO > 1 + + +def test_min_floor_pairs_rejects_negative(monkeypatch: pytest.MonkeyPatch) -> None: + db = _FakeDB() + with pytest.raises(ValueError, match="min_floor_pairs"): + _run(db, monkeypatch, min_floor_pairs=-1) + assert db.executed == [] + + +def test_min_floor_pairs_rejects_bool(monkeypatch: pytest.MonkeyPatch) -> None: + db = _FakeDB() + with pytest.raises(ValueError, match="min_floor_pairs"): + _run(db, monkeypatch, min_floor_pairs=True) + assert db.executed == [] + + +def test_floor_drop_ratio_rejects_below_one(monkeypatch: pytest.MonkeyPatch) -> None: + db = _FakeDB() + with pytest.raises(ValueError, match="floor_drop_ratio"): + _run(db, monkeypatch, floor_drop_ratio=0.5) + assert db.executed == [] + + +def test_floor_drop_ratio_rejects_bool(monkeypatch: pytest.MonkeyPatch) -> None: + db = _FakeDB() + with pytest.raises(ValueError, match="floor_drop_ratio"): + _run(db, monkeypatch, floor_drop_ratio=True) + assert db.executed == [] + + +def test_previous_floor_n_pairs_sql_compares_by_scrape_runs_source() -> None: + """Сравнение идёт по scrape_runs.source ТЕКУЩЕГО прогона (имя расписания), а + НЕ по listing_source -- см. модульный докстринг у DEFAULT_MIN_FLOOR_PAIRS.""" + sql = str(task_mod._PREVIOUS_FLOOR_N_PAIRS_SQL.text) + assert "SELECT source FROM scrape_runs WHERE id = :run_id" in sql + assert ":listing_source" not in sql + assert "status = 'done'" in sql + assert "id != :run_id" in sql + assert not re.search(r":\w+::", sql) + + +def test_handler_wires_floor_degradation_params_from_schedule() -> None: + """Читаем исходник файлом (как соседние test_handler_wires_* в этом сьюте): + product_handlers тянет scraper_kit, которого в юнит-окружении может не быть.""" + from pathlib import Path + + handlers = Path(__file__).resolve().parents[1] / "app" / "services" / "product_handlers.py" + src = handlers.read_text("utf-8") + job = src.split("async def _job_deactivate_stale")[1].split("\nasync def ")[0] + assert 'params.get("min_floor_pairs", DEFAULT_MIN_FLOOR_PAIRS)' in job + assert "min_floor_pairs=min_floor_pairs" in job + assert 'params.get("floor_drop_ratio", DEFAULT_FLOOR_DROP_RATIO)' in job + assert "floor_drop_ratio=floor_drop_ratio" in job diff --git a/tradein-mvp/uv.lock b/tradein-mvp/uv.lock index 7d3380a4..e58985d0 100644 --- a/tradein-mvp/uv.lock +++ b/tradein-mvp/uv.lock @@ -1490,7 +1490,7 @@ dependencies = [ [package.metadata] requires-dist = [ - { name = "curl-cffi", specifier = ">=0.7.0" }, + { name = "curl-cffi", specifier = ">=0.15.0" }, { name = "httpx", extras = ["socks"], specifier = ">=0.27.0" }, { name = "psycopg", extras = ["binary"], specifier = ">=3.2.0" }, { name = "pydantic", specifier = ">=2.7.0" }, @@ -1712,7 +1712,7 @@ dev = [ [package.metadata] requires-dist = [ { name = "bcrypt", specifier = ">=4.2.0" }, - { name = "curl-cffi", specifier = ">=0.7.0" }, + { name = "curl-cffi", specifier = ">=0.15.0" }, { name = "fastapi", specifier = ">=0.115.0" }, { name = "geoalchemy2", specifier = ">=0.15.0" }, { name = "httpx", extras = ["socks"], specifier = ">=0.27.0" },