feat(tradein/deactivate-stale): страховочные рельсы объёма снятия
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / changes (pull_request) Successful in 10s
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m34s
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / changes (pull_request) Successful in 10s
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m34s
Продолжение #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
This commit is contained in:
parent
1c16ea45a3
commit
22ccc1050c
5 changed files with 976 additions and 15 deletions
|
|
@ -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,
|
||||
),
|
||||
)
|
||||
|
||||
|
|
|
|||
|
|
@ -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`
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
@ -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
|
||||
4
tradein-mvp/uv.lock
generated
4
tradein-mvp/uv.lock
generated
|
|
@ -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" },
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue