feat(tradein/deactivate-stale): страховочные рельсы объёма снятия (PR-B) #3066
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.core.config import settings as _settings
|
||||||
from app.tasks.deactivate_stale_avito import (
|
from app.tasks.deactivate_stale_avito import (
|
||||||
CAP_MULT,
|
CAP_MULT,
|
||||||
|
DEFAULT_FLOOR_DROP_RATIO,
|
||||||
|
DEFAULT_MAX_DEACTIVATED,
|
||||||
DEFAULT_MIN_CONFIRMATIONS,
|
DEFAULT_MIN_CONFIRMATIONS,
|
||||||
|
DEFAULT_MIN_FLOOR_PAIRS,
|
||||||
DEFAULT_REVISIT_FLOOR_QUANTILE,
|
DEFAULT_REVISIT_FLOOR_QUANTILE,
|
||||||
deactivate_stale_listings,
|
deactivate_stale_listings,
|
||||||
)
|
)
|
||||||
|
|
@ -247,6 +250,14 @@ async def _job_deactivate_stale(
|
||||||
# относительно своего ttl_days переопределяет его через default_params (ключ
|
# относительно своего ttl_days переопределяет его через default_params (ключ
|
||||||
# "cap_mult"), не трогая дефолт для остальных источников.
|
# "cap_mult"), не трогая дефолт для остальных источников.
|
||||||
cap_mult: float = params.get("cap_mult", 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()
|
loop = asyncio.get_event_loop()
|
||||||
await loop.run_in_executor(
|
await loop.run_in_executor(
|
||||||
|
|
@ -262,6 +273,9 @@ async def _job_deactivate_stale(
|
||||||
revisit_floor_quantile=revisit_floor_quantile,
|
revisit_floor_quantile=revisit_floor_quantile,
|
||||||
null_segment_only=null_segment_only,
|
null_segment_only=null_segment_only,
|
||||||
cap_mult=cap_mult,
|
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
|
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:
|
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.
|
"""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
|
_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(
|
def deactivate_stale_listings(
|
||||||
db: Session,
|
db: Session,
|
||||||
run_id: int,
|
run_id: int,
|
||||||
|
|
@ -526,6 +697,9 @@ def deactivate_stale_listings(
|
||||||
revisit_floor_quantile: float = 0.0,
|
revisit_floor_quantile: float = 0.0,
|
||||||
null_segment_only: bool = False,
|
null_segment_only: bool = False,
|
||||||
cap_mult: float = CAP_MULT,
|
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]:
|
) -> dict[str, int]:
|
||||||
"""Пометить is_active=false объявления, чья свежесть старше ttl_days дней.
|
"""Пометить is_active=false объявления, чья свежесть старше ttl_days дней.
|
||||||
|
|
||||||
|
|
@ -568,20 +742,52 @@ def deactivate_stale_listings(
|
||||||
источники. Итоговый потолок = ttl_days * cap_mult. Применяется и к
|
источники. Итоговый потолок = ttl_days * cap_mult. Применяется и к
|
||||||
null_segment_only-джобе, но там гейт/пол выключены (см. выше), так что
|
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).
|
Sync (вызывается scheduler-триггером в executor, как snapshot_listing_sources).
|
||||||
Один statement в транзакции: UPDATE флага + снимок 'stale' в listings_snapshots
|
Один statement в транзакции: UPDATE флага + снимок 'stale' в listings_snapshots
|
||||||
(data-modifying CTE, #2674). Финализирует scrape_runs (mark_done / mark_failed).
|
(data-modifying CTE, #2674). Финализирует scrape_runs (mark_done / mark_failed).
|
||||||
|
|
||||||
Returns {"deactivated": N} -- количество обновлённых строк (1:1 со снимками).
|
Returns {"deactivated": N} -- количество обновлённых строк (1:1 со снимками).
|
||||||
Если гейт не пропустил прогон: {"deactivated": 0, "confirmations": N,
|
Если гейт здоровья не пропустил прогон: {"deactivated": 0, "confirmations": N,
|
||||||
"skipped_unhealthy": 1} и НИ ОДНА строка не тронута. Если пол переобхода включён
|
"skipped_unhealthy": 1} и НИ ОДНА строка не тронута. Если пол переобхода включён
|
||||||
(revisit_floor_quantile > 0): дополнительно {"floor_n_pairs": N} -- размер
|
(revisit_floor_quantile > 0): дополнительно {"floor_n_pairs": N} -- размер
|
||||||
выборки, из которой percentile_disc посчитал квантиль (наблюдательность, PR-A
|
выборки, из которой percentile_disc посчитал квантиль (PR-A), и (если найден
|
||||||
#2659 продолжение; сейчас ничего не гейтит). Если пол при этом реально поднял
|
предыдущий успешный прогон того же расписания) {"floor_n_pairs_previous": N}.
|
||||||
TTL: дополнительно {"revisit_floor_days": N, "ttl_days_effective": N}. Если пол
|
Если гейт деградации пола не пропустил прогон (PR-B): {"skipped_floor_degraded":
|
||||||
упёрся в потолок cap_mult: дополнительно {"ttl_floor_capped": 1,
|
1} и НИ ОДНА строка не тронута -- ни UPDATE, ни preflight-потолок ниже не
|
||||||
"ttl_days_floor_raw": N} -- N это то, во что пол поднял бы TTL БЕЗ потолка.
|
исполнялись. Если пол при этом реально поднял 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:
|
Raises:
|
||||||
ValueError: если staleness_column не входит в whitelist, ИЛИ ttl_days <= 0
|
ValueError: если staleness_column не входит в whitelist, ИЛИ ttl_days <= 0
|
||||||
|
|
@ -605,7 +811,11 @@ def deactivate_stale_listings(
|
||||||
которой оба guard'а вообще написаны, поэтому bool отклоняется явной
|
которой оба guard'а вообще написаны, поэтому bool отклоняется явной
|
||||||
type-проверкой ДО числового сравнения для обоих параметров), ЛИБО если
|
type-проверкой ДО числового сравнения для обоих параметров), ЛИБО если
|
||||||
заданы одновременно null_segment_only=True и segments (взаимоисключающие
|
заданы одновременно 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}
|
counters: dict[str, int] = {"deactivated": 0}
|
||||||
try:
|
try:
|
||||||
|
|
@ -631,6 +841,26 @@ def deactivate_stale_listings(
|
||||||
if cap_mult < 1:
|
if cap_mult < 1:
|
||||||
raise ValueError(f"cap_mult must be >= 1, got {cap_mult!r}")
|
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: только после неё имя колонки
|
# Whitelist-проверка ДО построения/выполнения SQL: только после неё имя колонки
|
||||||
# интерполируется f-string'ом. Значения по-прежнему идут через param-binding.
|
# интерполируется f-string'ом. Значения по-прежнему идут через param-binding.
|
||||||
# Внутри try -> невалидная колонка финализирует run как failed (mark_failed),
|
# Внутри try -> невалидная колонка финализирует run как failed (mark_failed),
|
||||||
|
|
@ -711,12 +941,13 @@ def deactivate_stale_listings(
|
||||||
),
|
),
|
||||||
floor_params,
|
floor_params,
|
||||||
).scalar()
|
).scalar()
|
||||||
# Наблюдательность (PR-A, #2659 продолжение) -- та же выборка, что и
|
# Размер выборки, из которой percentile_disc выше берёт квантиль (та же
|
||||||
# percentile_disc выше (общая _revisit_floor_from_where_sql). Пока ничего
|
# общая _revisit_floor_from_where_sql). Появилось в PR-A как чистое
|
||||||
# не гейтит (это PR-B), но n_pairs обязан попасть в counters ДО того, как
|
# наблюдение; с PR-B (#2659 продолжение) гейтит деградацию -- см. блок
|
||||||
# схлопнувшуюся выборку станет видно только по повторению прод-инцидента.
|
# ниже. n_pairs обязан попасть в counters ДО того, как схлопнувшуюся
|
||||||
# Считается ВСЕГДА при включённом полу, в том числе когда floor_days
|
# выборку станет видно только по повторению прод-инцидента. Считается
|
||||||
# получится NULL -- тогда n_pairs=0 и объясняет, почему пол не посчитался.
|
# ВСЕГДА при включённом полу, в том числе когда floor_days получится
|
||||||
|
# NULL -- тогда n_pairs=0 и объясняет, почему пол не посчитался.
|
||||||
floor_n_pairs = db.execute(
|
floor_n_pairs = db.execute(
|
||||||
_build_revisit_floor_pairs_count_sql(
|
_build_revisit_floor_pairs_count_sql(
|
||||||
staleness_column,
|
staleness_column,
|
||||||
|
|
@ -726,6 +957,44 @@ def deactivate_stale_listings(
|
||||||
floor_params,
|
floor_params,
|
||||||
).scalar()
|
).scalar()
|
||||||
counters["floor_n_pairs"] = int(floor_n_pairs or 0)
|
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 = истории снимков за окно нет вовсе (свежая БД, дыра в снимках).
|
# NULL = истории снимков за окно нет вовсе (свежая БД, дыра в снимках).
|
||||||
# Тогда пола нет и TTL остаётся как задан: выдумывать пол не из чего.
|
# Тогда пола нет и TTL остаётся как задан: выдумывать пол не из чего.
|
||||||
if floor_days is not None:
|
if floor_days is not None:
|
||||||
|
|
@ -779,6 +1048,56 @@ def deactivate_stale_listings(
|
||||||
ttl_days,
|
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_segment_only -> IS NULL, отдельный явный предикат (ANY(:segments)
|
||||||
# никогда не матчит NULL). segments is None -> все сегменты (поведение avito).
|
# никогда не матчит NULL). segments is None -> все сегменты (поведение avito).
|
||||||
# segments=[...] -> только перечисленные сегменты. Используем `is not None`
|
# 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]
|
[package.metadata]
|
||||||
requires-dist = [
|
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 = "httpx", extras = ["socks"], specifier = ">=0.27.0" },
|
||||||
{ name = "psycopg", extras = ["binary"], specifier = ">=3.2.0" },
|
{ name = "psycopg", extras = ["binary"], specifier = ">=3.2.0" },
|
||||||
{ name = "pydantic", specifier = ">=2.7.0" },
|
{ name = "pydantic", specifier = ">=2.7.0" },
|
||||||
|
|
@ -1712,7 +1712,7 @@ dev = [
|
||||||
[package.metadata]
|
[package.metadata]
|
||||||
requires-dist = [
|
requires-dist = [
|
||||||
{ name = "bcrypt", specifier = ">=4.2.0" },
|
{ 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 = "fastapi", specifier = ">=0.115.0" },
|
||||||
{ name = "geoalchemy2", specifier = ">=0.15.0" },
|
{ name = "geoalchemy2", specifier = ">=0.15.0" },
|
||||||
{ name = "httpx", extras = ["socks"], specifier = ">=0.27.0" },
|
{ name = "httpx", extras = ["socks"], specifier = ">=0.27.0" },
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue