fix(tradein/sber): сторож мерит отставание загрузки, а не календарь (#2846)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 9s
CI / frontend-tests (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m24s
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 9s
CI / frontend-tests (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m24s
Порог свежести якоря был недостижим по построению. period_month — метка ПЕРВОГО
числа месяца, поэтому возрасту ≥30 уже на закрытии месяца; плюс лаг публикации
источника. За 31 сутки прямых наблюдений монитора (07-13…08-12, scrape_runs.counters)
возраст лежал в 46..76 и ни разу не опускался ниже 46 — при пороге оценщика 35.
Сторож был истинным 100% времени с рождения таблицы и нёс ноль бит: при живом
источнике и при мёртвом загрузчике он писал одно и то же.
Второй порог (монитор, 60 = 35 + запас 25) не лучше: он лежит ВНУТРИ рабочего
диапазона. Обещание миграции 212 («такт 7 ⇒ потолок возраста 53 < 60») прод
ОПРОВЕРГ — 2026-08-12 возраст 72 при полном прогоне загрузки 08-06; двенадцатые
сутки подряд ERROR при исправной загрузке. Потолок 53 держался бы, только если бы
источник публиковал строго помесячно.
Что теперь. Загрузчик тянет ВСЮ серию (limit=1000&offset=0), поэтому после прогона
с errors=0 AND upserted>0 наш max(period_month) равен максимуму источника ПО
ПОСТРОЕНИЮ. Значит вопрос «отстали ли мы» = «давно ли был последний ЗАВЕДОМО ПОЛНЫЙ
прогон», и он не зависит от возраста периода. Порог — 2 такта самой загрузки,
читается из scrape_schedules.default_params.interval_days, то есть из той же строки,
по которой планировщик считает next_run_at: разъехаться с тактом он не может.
status='done' за успех не считается — прогон id=37 имеет done при {errors: 9,
upserted: 0}. Табло спрашивается тем же порядком, что у оценщика
(SBER_COEFF_DASHBOARDS), потому что max() по таблице маскирует отставшее табло:
real_estate_deals 2026-06, dinamika-tsen-obyavlenii 2026-05.
Два порога сведены удалением: settings.sber_index_max_age_days и per-estimate
warning в estimator убраны, свежесть считает ровно одно место.
fetched_at больше не переписывается апсертом. Забор идёт всей серией, поэтому
fetched_at = now() в DO UPDATE ставил одну метку всем 639 строкам, включая период
2017-01 — как признак свежести колонка была пуста. Теперь она означает «когда
впервые увидели период», т.е. такт публикации источника станет измеримым.
Ретроспективу это не возвращает: у уже лежащих строк метка 2026-08-06 и останется.
Refs #2846
This commit is contained in:
parent
e59a102b16
commit
b64e824e17
7 changed files with 541 additions and 263 deletions
|
|
@ -599,9 +599,13 @@ class Settings(BaseSettings):
|
|||
cian_valuation_max_rub: float = 500_000_000
|
||||
|
||||
# ── #audit-5: data-age guards ─────────────────────────────────────────────
|
||||
# sber_index_max_age_days: максимальный допустимый возраст последнего месяца
|
||||
# СберИндекс-серии (дней). Если latest месяц старее — логируем warning.
|
||||
sber_index_max_age_days: int = 35
|
||||
# #2846: sber_index_max_age_days УДАЛЁН (был 35). Порог недостижим по построению
|
||||
# (period_month — метка первого числа + лаг публикации источника ⇒ пол 46 суток),
|
||||
# guard был истинным 100% времени. Свежесть СберИндекса теперь считает ровно одно
|
||||
# место — tasks/sber_freshness_monitor, и считает по отставанию ЗАГРУЗКИ, а порог
|
||||
# берёт из такта самой загрузки (scrape_schedules.default_params.interval_days),
|
||||
# так что второму порогу тут больше неоткуда взяться и не с чем разъезжаться.
|
||||
# extra="ignore" в model_config защищает от startup-краха на leftover env var.
|
||||
# avito_imv_thin_market_threshold: если market_count < порога — IMV-оценка
|
||||
# на тонком рынке (thin_market=True в AvitoImvSummary) + warning.
|
||||
avito_imv_thin_market_threshold: int = 10
|
||||
|
|
|
|||
|
|
@ -1565,7 +1565,16 @@ def _load_sber_index_series(db: Session, *, region: str) -> dict[date, float]:
|
|||
"""#794: monthly {period_month: index_value} for region from sber_price_index.
|
||||
|
||||
Tries SBER_COEFF_DASHBOARDS in order; returns first non-empty series. {} on any error.
|
||||
#audit-5a: если latest месяц серии старее sber_index_max_age_days → warning.
|
||||
|
||||
#2846: per-estimate guard свежести отсюда УБРАН. Он сравнивал возраст latest
|
||||
периода с settings.sber_index_max_age_days=35, а такой возраст недостижим по
|
||||
построению: period_month — метка ПЕРВОГО числа месяца (≥30 суток уже на
|
||||
закрытии месяца) плюс лаг публикации источника; на проде за 31 сутки прямых
|
||||
наблюдений возраст не опускался ниже 46. Guard был истинным 100% времени —
|
||||
нулевой сигнал в per-estimate логе, который вдобавок не долетал до GlitchTip
|
||||
(event_level=ERROR). Свежесть теперь мерит ОДНО место — tasks/sber_freshness_monitor,
|
||||
и мерит отставание ЗАГРУЗКИ (последний полный прогон vs её собственный такт),
|
||||
а не календарь.
|
||||
"""
|
||||
for dash in SBER_COEFF_DASHBOARDS:
|
||||
try:
|
||||
|
|
@ -1593,21 +1602,6 @@ def _load_sber_index_series(db: Session, *, region: str) -> dict[date, float]:
|
|||
series = {r["period_month"]: float(r["index_value_rub_m2"]) for r in rows}
|
||||
if not series:
|
||||
continue
|
||||
# #audit-5a: data-age guard — предупреждаем о stale СберИндексе.
|
||||
latest = max(series)
|
||||
today = datetime.now(tz=UTC).date()
|
||||
age_days = (today - latest).days
|
||||
if age_days > settings.sber_index_max_age_days:
|
||||
logger.warning(
|
||||
"sber_index stale #audit-5a: latest=%s age=%d days"
|
||||
" (> sber_index_max_age_days=%d) region=%s dash=%s"
|
||||
" — time-adjustment may be outdated",
|
||||
latest.isoformat(),
|
||||
age_days,
|
||||
settings.sber_index_max_age_days,
|
||||
region,
|
||||
dash,
|
||||
)
|
||||
return series
|
||||
return {}
|
||||
|
||||
|
|
|
|||
|
|
@ -303,6 +303,18 @@ def _upsert_rows_sync(db: Session, rows_to_upsert: list[tuple[str, date, str, st
|
|||
|
||||
#1348: blocking psycopg work — must run via asyncio.to_thread, never directly
|
||||
on the event loop. Idempotent ON CONFLICT(city, period_month, dashboard).
|
||||
|
||||
#2846: `fetched_at` НЕ переписывается при конфликте. Забор идёт ВСЕЙ серией
|
||||
(limit=1000&offset=0, отсечки по периоду нет), поэтому `fetched_at = now()` в
|
||||
DO UPDATE ставил одну и ту же метку всем строкам ряда — на проде все 639 строк
|
||||
несли время последнего прогона, включая период 2017-01. Как признак свежести
|
||||
колонка была пуста. Теперь она означает «когда мы ВПЕРВЫЕ увидели этот период»,
|
||||
то есть по ней измеряется ТАКТ ПУБЛИКАЦИИ источника (min(fetched_at) по новым
|
||||
периодам). Ретроспективу это не возвращает: у 639 уже лежащих строк метка
|
||||
2026-08-06 и она останется — такт публикации до этого PR невосстановим.
|
||||
Времени последней ЗАГРУЗКИ колонка больше не хранит; оно и не нужно —
|
||||
scrape_runs(source='sber_index_pull') хранит его точнее (с errors/upserted),
|
||||
и именно оттуда его берёт tasks/sber_freshness_monitor.
|
||||
"""
|
||||
for city_label, period_month, segment, dash, value in rows_to_upsert:
|
||||
db.execute(
|
||||
|
|
@ -323,8 +335,8 @@ def _upsert_rows_sync(db: Session, rows_to_upsert: list[tuple[str, date, str, st
|
|||
ON CONFLICT (city, period_month, dashboard)
|
||||
DO UPDATE SET
|
||||
index_value_rub_m2 = EXCLUDED.index_value_rub_m2,
|
||||
segment = EXCLUDED.segment,
|
||||
fetched_at = now()
|
||||
segment = EXCLUDED.segment
|
||||
-- fetched_at НЕ трогаем (#2846): она = «впервые увидели период».
|
||||
"""
|
||||
),
|
||||
{
|
||||
|
|
|
|||
|
|
@ -1,56 +1,68 @@
|
|||
"""Мониторинг свежести ДАННЫХ СберИндекса (не статуса джобы) — audit п.1.
|
||||
"""Монитор ОТСТАВАНИЯ ЗАГРУЗКИ СберИндекса (не календарного возраста периода).
|
||||
|
||||
Проблема аудита: estimator._load_sber_index_series (#794/#audit-5a) применяет
|
||||
СберИндекс time-adjustment к ДКП-сделкам и лишь ЛОГИРУЕТ per-estimate warning,
|
||||
когда latest месяц серии старее settings.sber_index_max_age_days (35д). Джоба
|
||||
`sber_index_pull` крутится ежемесячно (enabled), а источник СберИндекса публикует
|
||||
данные с лагом ~1-2 месяца, поэтому `sber_price_index.period_month` дрейфит
|
||||
(на 2026-07-12 latest=2026-05-01, ~72д). Это НЕ silent failure, но staleness
|
||||
видна только в debug-подобном per-estimate warning'е, тонущем в логах оценок.
|
||||
ЧТО БЫЛО НЕ ТАК (замер на проде 2026-08-12, #2846).
|
||||
|
||||
Этот монитор смотрит на `max(period_month)` вторичного сегмента по региону и
|
||||
поднимает per-day ERROR-алерт, когда данные устарели СВЕРХ допустимого лага
|
||||
публикации — так ops видит дрейф на MONITOR-частоте, а не по крупицам в логах.
|
||||
Монитор мерил `now() - max(period_month)` и алертил при возрасте > 60 суток
|
||||
(sber_index_max_age_days 35 + lag_allowance 25). Такой возраст НЕДОСТИЖИМО МАЛ по
|
||||
построению: `period_month` — метка ПЕРВОГО числа месяца, поэтому на закрытии месяца
|
||||
возрасту уже ≥30; плюс собственный лаг публикации источника. За 31 сутки прямых
|
||||
наблюдений монитора (07-13 … 08-12, scrape_runs.counters) возраст лежал в 46..76 и
|
||||
НИ РАЗУ не опускался ниже 46. Порог 35 у оценщика был истинным 100% времени — ноль бит.
|
||||
|
||||
#2674 — почему ERROR, а не WARNING. В контейнере скрапера GlitchTip поднят с
|
||||
LoggingIntegration(event_level=ERROR) (scheduler_main.py), поэтому WARNING
|
||||
событием НЕ становится вообще. Бенчмарк цен участвует в сверке наших медиан, его
|
||||
застой — сбой, а не наблюдение. Сосед по конструкции (deals_freshness_monitor)
|
||||
писал ERROR с самого начала — расходилась только эта джоба.
|
||||
Порог 60 у монитора не лучше: он лежит ВНУТРИ рабочего диапазона, поэтому монитор
|
||||
мерил не источник, а нашу же пилу. Миграция 212 (такт 28 → 7) обещала потолок
|
||||
возраста ≈46+7=53 < 60. Прод это ОПРОВЕРГ: 2026-08-12 возраст 72 при ПОЛНОМ прогоне
|
||||
загрузки шестидневной давности (08-06, errors=0, upserted=639) — источник просто не
|
||||
опубликовал июль. Двенадцать суток подряд (08-01 … 08-12) монитор писал ERROR при
|
||||
исправной загрузке. Потолок 53 держится, только если источник публикует строго
|
||||
помесячно; он не публикует.
|
||||
|
||||
ВАЖНО про «9 срабатываний» из #2674 (ревью PR #2681, прод-разбор всех 24 прогонов
|
||||
монитора 2026-08-06). Эти девять НЕ были застоем бенчмарка — это была ПИЛА нашего
|
||||
собственного такта загрузки:
|
||||
13-16.07 alert=1 age 73..76 latest=май 01-05.08 alert=1 age 61..65
|
||||
17.07 alert=0 age 46 latest=июнь (день загрузки)
|
||||
Загрузка ходила раз в 28 дней и приносила период на месяц новее, возраст же
|
||||
считается от ПЕРВОГО числа покрытого месяца → пол ~46 в момент загрузки, потолок
|
||||
46+28=74, порог 60 ВНУТРИ диапазона, тревога 14 суток из 28 каждый цикл. Поднимать
|
||||
такое до ERROR без починки такта значило бы завести ежедневное ложное событие на
|
||||
две недели в месяц. Поэтому миграция 212 перевела sber_index_pull на НЕДЕЛЬНЫЙ
|
||||
такт: потолок возраста ≈ пол+7 ≈ 53 при пороге 60, тревога снова означает
|
||||
«источник/загрузка встали», а не «мы давно не ходили».
|
||||
ЧТО МЕРИМ ТЕПЕРЬ. Загрузчик тянет ВСЮ серию (limit=1000&offset=0, отсечки по периоду
|
||||
нет), поэтому после прогона с errors=0 AND upserted>0 наш max(period_month) РАВЕН
|
||||
максимуму источника ПО ПОСТРОЕНИЮ. Значит вопрос «отстали ли мы» — это вопрос
|
||||
«давно ли был последний ЗАВЕДОМО ПОЛНЫЙ прогон», и он не зависит от возраста периода:
|
||||
|
||||
Порог алерта (документирование выбора):
|
||||
Per-estimate guard (estimator): age > settings.sber_index_max_age_days (35д).
|
||||
Монитор: age > sber_index_max_age_days + lag_allowance.
|
||||
lag_allowance (DEFAULT_LAG_ALLOWANCE_DAYS=25) — запас на ИНХЕРЕНТНЫЙ лаг
|
||||
публикации СберИндекса: источник отстаёт на 1-2 месяца, а period_month — лейбл
|
||||
ПЕРВОГО числа месяца, поэтому даже свежайшая загрузка даёт возраст ~46 суток.
|
||||
Итог: 35 + 25 = 60д. При недельном такте (миграция 212) рабочий диапазон возраста
|
||||
~46..53 — до порога остаётся ~7 суток запаса: один пропущенный недельный цикл
|
||||
поглощается, два подряд дают тревогу. Порог НЕ должен снова оказаться внутри
|
||||
рабочего диапазона — если такт загрузки будут менять, пересчитай потолок
|
||||
(пол + interval_days) и сверь с 60.
|
||||
последний полный прогон свежий → наш max == max источника → источник не публиковал,
|
||||
молчание ПРАВИЛЬНОЕ (возраст = лаг источника);
|
||||
последний полный прогон старый → мы не забрали → тревога про ЗАГРУЗЧИК.
|
||||
|
||||
Задача синхронная (DB-only, один SELECT max(period_month)) — запускается
|
||||
kit-scheduler'ом через product_handlers._job_sber_freshness_monitor в
|
||||
run_in_executor, по образцу deals_freshness_monitor. Вердикт вычисляет ЧИСТАЯ
|
||||
функция evaluate_sber_freshness() (frozen-now тестируется без БД).
|
||||
ЛОВУШКА: `status='done'` НЕ означает успех — прогон id=37 (2026-05-31) имеет
|
||||
{errors: 9, upserted: 0} и статус done. Успех = errors=0 AND upserted>0 (все 9 серий
|
||||
3 табло × 3 региона прошли: errors — счётчик по всему прогону).
|
||||
|
||||
Прогон НЕ помечается failed при алерте (это МОНИТОР, а не сбой джобы) — ERROR-записи
|
||||
достаточно. mark_failed только если sber_price_index недоступна/пуста (нечего
|
||||
оценивать).
|
||||
ПОРОГ — не круглое число, а такт самой загрузки: `scrape_schedules.default_params
|
||||
.interval_days` для sber_index_pull, ЧИТАЕТСЯ ИЗ ТОЙ ЖЕ СТРОКИ, по которой планировщик
|
||||
запускает прогон (orchestration/scheduler.py::_defer_next_run_at). Разъехаться с
|
||||
тактом порог не может: поменяли такт — порог поехал следом. Тревога после
|
||||
MISSED_PULL_CYCLES=2 пропущенных тактов: один пропуск (сдвиг окна, разовый сбой сети)
|
||||
поглощается, два подряд означают, что загрузка встала. При нынешнем такте 7 это 14
|
||||
суток; на прод-истории такое состояние ДОСТИЖИМО — разрывы между полными прогонами
|
||||
были 14.8 и 20 суток (05-31→06-15 и 07-17→08-06).
|
||||
|
||||
ПО ТАБЛО, А НЕ ПО max() ВСЕЙ ТАБЛИЦЫ. Оценщик берёт ПЕРВОЕ НЕПУСТОЕ табло из
|
||||
estimator.SBER_COEFF_DASHBOARDS; у real_estate_deals latest=2026-06, у
|
||||
dinamika-tsen-obyavlenii — 2026-05 (на 2026-08-12). max() по таблице маскирует
|
||||
отставшее табло, поэтому монитор идёт тем же порядком, что и оценщик, и берёт ту же
|
||||
серию — список импортируется из estimator, дублировать его тут нельзя.
|
||||
|
||||
ЧЕГО ЭТОТ МОНИТОР НЕ ЛОВИТ (осознанно, #2846). Если источник ЗАМОЛЧИТ НАВСЕГДА, а
|
||||
загрузка останется исправной — монитор промолчит: по нашим данным «источник не
|
||||
публиковал 2 месяца» неотличимо от «источник публикует раз в 2 месяца». Такт
|
||||
публикации источника ретроспективно невосстановим — его затёр апсерт
|
||||
(sber_index.py ставил fetched_at=now() всем строкам серии). С этого PR fetched_at
|
||||
не переписывается при конфликте и означает «когда мы ВПЕРВЫЕ увидели этот период»,
|
||||
т.е. такт публикации станет измеримым; вернуться к вопросу порога «источник встал»
|
||||
имеет смысл после 3 наблюдённых публикаций (ориентир — ноябрь 2026).
|
||||
|
||||
ERROR, а не WARNING (#2674): в контейнере скрапера GlitchTip поднят с
|
||||
LoggingIntegration(event_level=ERROR), WARNING событием не становится вообще.
|
||||
|
||||
Задача синхронная (DB-only) — запускается kit-scheduler'ом через
|
||||
product_handlers._job_sber_freshness_monitor в run_in_executor. Вердикт считает
|
||||
ЧИСТАЯ функция evaluate_sber_freshness() (frozen-now, тестируется без БД).
|
||||
|
||||
Прогон НЕ помечается failed при алерте (это МОНИТОР, а не сбой джобы). mark_failed
|
||||
только если у оценщика вообще нет серии (нечего оценивать).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
|
@ -62,153 +74,239 @@ from datetime import UTC, date, datetime
|
|||
from sqlalchemy import text
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from app.core.config import settings
|
||||
from app.services import scrape_runs as runs_mod
|
||||
from app.services.estimator import SBER_COEFF_DASHBOARDS, SBER_TIME_ADJUST_REGION
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
__all__ = [
|
||||
"DEFAULT_LAG_ALLOWANCE_DAYS",
|
||||
"DEFAULT_PULL_INTERVAL_DAYS",
|
||||
"MISSED_PULL_CYCLES",
|
||||
"SBER_FRESHNESS_PULL_SOURCE",
|
||||
"SberFreshnessVerdict",
|
||||
"check_sber_freshness",
|
||||
"evaluate_sber_freshness",
|
||||
]
|
||||
|
||||
# Запас на инхерентный лаг публикации СберИндекса (дней) СВЕРХ per-estimate
|
||||
# guard'а settings.sber_index_max_age_days. Читается из default_params.lag_allowance_days.
|
||||
DEFAULT_LAG_ALLOWANCE_DAYS = 25
|
||||
# Джоба-загрузчик, чей такт и успешность мы и мониторим.
|
||||
SBER_FRESHNESS_PULL_SOURCE = "sber_index_pull"
|
||||
|
||||
# Регион продукта (Trade-in — Свердловская область). Совпадает с city-значениями
|
||||
# sber_price_index для областного уровня.
|
||||
SBER_MONITOR_CITY = "Свердловская область"
|
||||
# Сколько тактов загрузки подряд можно пропустить до тревоги. 1 = разовый сбой/сдвиг
|
||||
# окна (поглощаем), 2 = загрузка встала (алерт).
|
||||
MISSED_PULL_CYCLES = 2
|
||||
|
||||
_LATEST_SBER_PERIOD_SQL = text("""
|
||||
# Фолбэк, если в scrape_schedules нет строки/ключа interval_days (миграция 212 ставит 7).
|
||||
DEFAULT_PULL_INTERVAL_DAYS = 7
|
||||
|
||||
_LATEST_PERIOD_SQL = text("""
|
||||
SELECT max(period_month) AS latest
|
||||
FROM sber_price_index
|
||||
WHERE city = CAST(:city AS text)
|
||||
AND dashboard = CAST(:dash AS text)
|
||||
-- #R2-H1: только вторичный рынок (эстиматор — вторичка); первичка
|
||||
-- (новостройки) = направленно неверная коррекция. Зеркалит фильтр
|
||||
-- estimator._load_sber_index_series.
|
||||
AND (segment IS NULL OR segment ILIKE '%вторичн%')
|
||||
""")
|
||||
|
||||
# Последний ЗАВЕДОМО ПОЛНЫЙ прогон загрузчика. status='done' сюда не входит намеренно:
|
||||
# прогон id=37 имеет done при {errors: 9, upserted: 0}. Сравнения — jsonb-ные, без
|
||||
# CAST(... AS int): counters других источников планировщик может отфильтровать позже
|
||||
# каста, а не раньше, и нечисловое значение уронило бы запрос. Для jsonb-чисел
|
||||
# оператор > численный.
|
||||
_LAST_COMPLETE_PULL_SQL = text("""
|
||||
SELECT max(finished_at) AS last_pull
|
||||
FROM scrape_runs
|
||||
WHERE source = CAST(:src AS text)
|
||||
AND counters @> CAST('{"errors": 0}' AS jsonb)
|
||||
AND counters -> 'upserted' > CAST('0' AS jsonb)
|
||||
""")
|
||||
|
||||
# Такт загрузки — из той же строки, по которой планировщик считает next_run_at.
|
||||
_PULL_INTERVAL_SQL = text("""
|
||||
SELECT default_params ->> 'interval_days' AS interval_days
|
||||
FROM scrape_schedules
|
||||
WHERE source = CAST(:src AS text)
|
||||
""")
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class SberFreshnessVerdict:
|
||||
"""Вердикт свежести СберИндекса по max(period_month)."""
|
||||
"""Вердикт: отстала ли ЗАГРУЗКА СберИндекса от собственного такта."""
|
||||
|
||||
latest_period: date
|
||||
age_days: int
|
||||
age_days: int # наблюдение (лаг публикации источника), НЕ критерий тревоги
|
||||
pull_lag_days: int # суток с последнего полного прогона; -1 = полных прогонов не было
|
||||
max_pull_lag_days: int # порог = MISSED_PULL_CYCLES × такт загрузки
|
||||
stale: bool
|
||||
|
||||
|
||||
def evaluate_sber_freshness(
|
||||
latest_period: date,
|
||||
now: datetime,
|
||||
max_age_days: int,
|
||||
*,
|
||||
last_complete_pull_at: datetime | None,
|
||||
pull_interval_days: int,
|
||||
) -> SberFreshnessVerdict:
|
||||
"""Чистая логика: устарел ли latest период СберИндекса.
|
||||
"""Чистая логика: отстала ли загрузка от собственного такта.
|
||||
|
||||
stale = age_days > max_age_days, где age_days = now.date() - latest_period.
|
||||
`max_age_days` — ПОЛНЫЙ порог монитора (per-estimate guard + lag_allowance),
|
||||
вычисляется вызывающим check_sber_freshness. Тестируется с frozen `now` без БД.
|
||||
stale = полных прогонов не было ВООБЩЕ, либо последний старше
|
||||
MISSED_PULL_CYCLES × pull_interval_days. Возраст периода считается и кладётся в
|
||||
вердикт как НАБЛЮДЕНИЕ, но на вердикт не влияет: после полного прогона наш
|
||||
max(period_month) равен максимуму источника по построению, и его возраст — это
|
||||
лаг ПУБЛИКАЦИИ, на который мы повлиять не можем.
|
||||
"""
|
||||
age_days = (now.date() - latest_period).days
|
||||
stale = age_days > max_age_days
|
||||
max_pull_lag_days = MISSED_PULL_CYCLES * pull_interval_days
|
||||
if last_complete_pull_at is None:
|
||||
return SberFreshnessVerdict(latest_period, age_days, -1, max_pull_lag_days, True)
|
||||
pull_lag_days = (now - last_complete_pull_at).days
|
||||
return SberFreshnessVerdict(
|
||||
latest_period=latest_period,
|
||||
age_days=age_days,
|
||||
stale=stale,
|
||||
pull_lag_days=pull_lag_days,
|
||||
max_pull_lag_days=max_pull_lag_days,
|
||||
stale=pull_lag_days > max_pull_lag_days,
|
||||
)
|
||||
|
||||
|
||||
def _load_estimator_dashboard(db: Session) -> tuple[str, date] | None:
|
||||
"""Табло, которое возьмёт оценщик, и его latest период.
|
||||
|
||||
Тот же порядок, что и estimator._load_sber_index_series: первое НЕПУСТОЕ табло
|
||||
из SBER_COEFF_DASHBOARDS. max() по всей таблице маскировал бы отставшее табло.
|
||||
"""
|
||||
for dash in SBER_COEFF_DASHBOARDS:
|
||||
row = db.execute(
|
||||
_LATEST_PERIOD_SQL, {"city": SBER_TIME_ADJUST_REGION, "dash": dash}
|
||||
).first()
|
||||
latest = row.latest if row is not None else None
|
||||
if latest is not None:
|
||||
return dash, latest
|
||||
return None
|
||||
|
||||
|
||||
def _pull_interval_days(db: Session) -> int:
|
||||
"""Такт загрузчика из scrape_schedules (фолбэк DEFAULT_PULL_INTERVAL_DAYS)."""
|
||||
row = db.execute(_PULL_INTERVAL_SQL, {"src": SBER_FRESHNESS_PULL_SOURCE}).first()
|
||||
raw = row.interval_days if row is not None else None
|
||||
try:
|
||||
return int(raw) if raw is not None else DEFAULT_PULL_INTERVAL_DAYS
|
||||
except (TypeError, ValueError):
|
||||
logger.warning(
|
||||
"sber freshness: interval_days=%r в scrape_schedules нечисловой — беру %d",
|
||||
raw,
|
||||
DEFAULT_PULL_INTERVAL_DAYS,
|
||||
)
|
||||
return DEFAULT_PULL_INTERVAL_DAYS
|
||||
|
||||
|
||||
def check_sber_freshness(
|
||||
db: Session,
|
||||
run_id: int,
|
||||
params: dict | None = None, # type: ignore[type-arg]
|
||||
now: datetime | None = None,
|
||||
) -> dict[str, int]:
|
||||
"""Проверить свежесть СберИндекса по max(period_month) и алертить при staleness.
|
||||
"""Проверить, не отстала ли загрузка СберИндекса, и алертить при отставании.
|
||||
|
||||
Sync (вызывается scheduler-триггером в executor, как check_deals_freshness).
|
||||
Читает один SELECT max(period_month) вторичного сегмента по региону, считает
|
||||
вердикт чистой функцией, логирует WARNING при stale (per-day surfacing для ops)
|
||||
и финализирует run.
|
||||
Читает: latest период табло оценщика, время последнего ПОЛНОГО прогона
|
||||
sber_index_pull, такт загрузки из scrape_schedules. Вердикт — чистой функцией.
|
||||
|
||||
Params (default_params jsonb):
|
||||
lag_allowance_days: int — запас сверх sber_index_max_age_days (default 25).
|
||||
`now` инъектируется в тестах (frozen); в проде — None → datetime.now(UTC).
|
||||
`params` больше ничего не настраивает: порог берётся из такта самой загрузки
|
||||
(унаследованный default_params.lag_allowance_days=25 монитора игнорируется —
|
||||
он кодировал мёртвый календарный порог). `now` инъектируется в тестах.
|
||||
|
||||
Returns counters {latest_year, latest_month, age_days, alert}.
|
||||
mark_failed только если sber_price_index пуста/недоступна (нечего оценивать);
|
||||
Returns counters {latest_year, latest_month, age_days, pull_lag_days,
|
||||
max_pull_lag_days, alert}.
|
||||
mark_failed только если у оценщика нет серии вообще (нечего оценивать);
|
||||
при алерте прогон помечается done (это монитор, не сбой джобы).
|
||||
"""
|
||||
params = params or {}
|
||||
now = now or datetime.now(UTC)
|
||||
counters: dict[str, int] = {
|
||||
"latest_year": 0,
|
||||
"latest_month": 0,
|
||||
"age_days": 0,
|
||||
"pull_lag_days": -1,
|
||||
"max_pull_lag_days": 0,
|
||||
"alert": 0,
|
||||
}
|
||||
try:
|
||||
runs_mod.update_heartbeat(db, run_id, counters)
|
||||
|
||||
row = db.execute(_LATEST_SBER_PERIOD_SQL, {"city": SBER_MONITOR_CITY}).first()
|
||||
latest: date | None = row.latest if row is not None else None
|
||||
if latest is None:
|
||||
found = _load_estimator_dashboard(db)
|
||||
if found is None:
|
||||
# ERROR (#2674): монитор не может выполнить свою работу вовсе — это сбой,
|
||||
# а не наблюдение. mark_failed ниже виден только стрик-алерту (3 подряд),
|
||||
# а монитор ходит раз в сутки — три дня молчания на пустом бенчмарке.
|
||||
logger.error(
|
||||
"sber freshness: sber_price_index пуст/недоступен для region=%s "
|
||||
"(вторичка) — оценить свежесть нельзя",
|
||||
SBER_MONITOR_CITY,
|
||||
"sber freshness: у оценщика нет серии — ни одно табло %s не даёт строк "
|
||||
"для region=%s (вторичка); оценить нечего",
|
||||
list(SBER_COEFF_DASHBOARDS),
|
||||
SBER_TIME_ADJUST_REGION,
|
||||
)
|
||||
runs_mod.mark_failed(db, run_id, "sber_price_index empty or unavailable", counters)
|
||||
return counters
|
||||
|
||||
lag_days = int(params.get("lag_allowance_days", DEFAULT_LAG_ALLOWANCE_DAYS))
|
||||
max_age_days = settings.sber_index_max_age_days + lag_days
|
||||
verdict = evaluate_sber_freshness(latest, now, max_age_days)
|
||||
dashboard, latest = found
|
||||
last_pull_row = db.execute(
|
||||
_LAST_COMPLETE_PULL_SQL, {"src": SBER_FRESHNESS_PULL_SOURCE}
|
||||
).first()
|
||||
last_complete_pull_at = last_pull_row.last_pull if last_pull_row is not None else None
|
||||
verdict = evaluate_sber_freshness(
|
||||
latest,
|
||||
now,
|
||||
last_complete_pull_at=last_complete_pull_at,
|
||||
pull_interval_days=_pull_interval_days(db),
|
||||
)
|
||||
|
||||
counters = {
|
||||
"latest_year": latest.year,
|
||||
"latest_month": latest.month,
|
||||
"age_days": verdict.age_days,
|
||||
"pull_lag_days": verdict.pull_lag_days,
|
||||
"max_pull_lag_days": verdict.max_pull_lag_days,
|
||||
"alert": int(verdict.stale),
|
||||
}
|
||||
|
||||
if verdict.stale:
|
||||
# ERROR (#2674): WARNING не долетает до GlitchTip (event_level=ERROR) —
|
||||
# 9 срабатываний на проде дали ноль событий. См. докстринг модуля.
|
||||
# ERROR (#2674): WARNING не долетает до GlitchTip (event_level=ERROR).
|
||||
logger.error(
|
||||
"sber freshness: max(period_month)=%s устарел на %d дней "
|
||||
"(> порога %d = sber_index_max_age_days %d + lag %d); "
|
||||
"СберИндекс time-adjustment ДКП-сделок мог отстать — "
|
||||
"проверь sber_index_pull и доступность новых периодов источника",
|
||||
"sber freshness: загрузка СберИндекса отстала — последний ПОЛНЫЙ прогон "
|
||||
"%s (%s суток назад, порог %d = %d такта × %d суток; "
|
||||
"status='done' с errors>0 за успех НЕ считается). "
|
||||
"Наш max(period_month)=%s (табло %s) мог разойтись с источником — "
|
||||
"проверь sber_index_pull: планировщик, сеть, /api/sowa 404",
|
||||
last_complete_pull_at.isoformat() if last_complete_pull_at else "НИ РАЗУ",
|
||||
verdict.pull_lag_days if verdict.pull_lag_days >= 0 else "∞",
|
||||
verdict.max_pull_lag_days,
|
||||
MISSED_PULL_CYCLES,
|
||||
verdict.max_pull_lag_days // MISSED_PULL_CYCLES,
|
||||
latest,
|
||||
verdict.age_days,
|
||||
max_age_days,
|
||||
settings.sber_index_max_age_days,
|
||||
lag_days,
|
||||
dashboard,
|
||||
)
|
||||
else:
|
||||
logger.info(
|
||||
"sber freshness: max(period_month)=%s свежий (age=%d дней ≤ порога %d) "
|
||||
"region=%s — алерта нет",
|
||||
"sber freshness: загрузка в такте — последний полный прогон %d суток назад "
|
||||
"(≤ порога %d). max(period_month)=%s (табло %s, возраст %d суток) равен "
|
||||
"максимуму источника по построению: возраст = лаг ПУБЛИКАЦИИ источника, "
|
||||
"не наше отставание — алерта нет",
|
||||
verdict.pull_lag_days,
|
||||
verdict.max_pull_lag_days,
|
||||
latest,
|
||||
dashboard,
|
||||
verdict.age_days,
|
||||
max_age_days,
|
||||
SBER_MONITOR_CITY,
|
||||
)
|
||||
|
||||
runs_mod.mark_done(db, run_id, counters)
|
||||
logger.info(
|
||||
"check_sber_freshness run_id=%d done: latest=%s alert=%d age_days=%d",
|
||||
"check_sber_freshness run_id=%d done: latest=%s dash=%s alert=%d "
|
||||
"pull_lag_days=%d age_days=%d",
|
||||
run_id,
|
||||
latest,
|
||||
dashboard,
|
||||
counters["alert"],
|
||||
counters["pull_lag_days"],
|
||||
counters["age_days"],
|
||||
)
|
||||
return counters
|
||||
|
|
|
|||
|
|
@ -102,14 +102,32 @@ def test_harness_itself_drops_warnings() -> None:
|
|||
|
||||
|
||||
class _FakeMonitorDB:
|
||||
"""Session-мок мониторов свежести: один SELECT max(...)."""
|
||||
"""Session-мок мониторов свежести.
|
||||
|
||||
def __init__(self, latest: date | None) -> None:
|
||||
Монитор сделок спрашивает только max(...). Монитор СберИндекса (#2846) спрашивает
|
||||
ещё время последнего ПОЛНОГО прогона загрузки и её такт — именно они, а не
|
||||
календарный возраст периода, решают, быть ли тревоге.
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
latest: date | None,
|
||||
last_pull: datetime | None = datetime(2026, 8, 6, 5, 0, tzinfo=UTC),
|
||||
interval_days: str = "7",
|
||||
) -> None:
|
||||
self._latest = latest
|
||||
self._last_pull = last_pull
|
||||
self._interval_days = interval_days
|
||||
|
||||
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
|
||||
sql = str(stmt)
|
||||
result = MagicMock()
|
||||
result.first.return_value = MagicMock(latest=self._latest)
|
||||
if "scrape_runs" in sql:
|
||||
result.first.return_value = MagicMock(last_pull=self._last_pull)
|
||||
elif "scrape_schedules" in sql:
|
||||
result.first.return_value = MagicMock(interval_days=self._interval_days)
|
||||
else:
|
||||
result.first.return_value = MagicMock(latest=self._latest)
|
||||
return result
|
||||
|
||||
def rollback(self) -> None:
|
||||
|
|
@ -122,10 +140,13 @@ def _patch_runs(monkeypatch: pytest.MonkeyPatch, module: Any) -> None:
|
|||
monkeypatch.setattr(module.runs_mod, "mark_failed", lambda *a, **k: None)
|
||||
|
||||
|
||||
def test_sber_staleness_becomes_event(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Прод-состояние (9 срабатываний, ноль событий): застой бенчмарка → событие."""
|
||||
def test_sber_pull_stall_becomes_event(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Загрузка встала (полный прогон 20 суток назад при такте 7) → событие.
|
||||
|
||||
20 суток — реальный разрыв прод-истории между полными прогонами 07-17 и 08-06.
|
||||
"""
|
||||
_patch_runs(monkeypatch, sber_mon)
|
||||
db = _FakeMonitorDB(date(2026, 6, 1))
|
||||
db = _FakeMonitorDB(date(2026, 6, 1), last_pull=datetime(2026, 7, 17, 5, 38, tzinfo=UTC))
|
||||
with glitchtip_events() as events:
|
||||
out = sber_mon.check_sber_freshness(
|
||||
db, # type: ignore[arg-type]
|
||||
|
|
@ -136,19 +157,24 @@ def test_sber_staleness_becomes_event(monkeypatch: pytest.MonkeyPatch) -> None:
|
|||
assert out["alert"] == 1
|
||||
assert any(
|
||||
"sber freshness" in t for t in event_texts(events)
|
||||
), "устаревание СберИндекса не стало событием — WARNING до GlitchTip не долетает"
|
||||
), "отставание загрузки не стало событием — WARNING до GlitchTip не долетает"
|
||||
|
||||
|
||||
def test_sber_fresh_index_stays_silent(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Свежие данные — ни одного события (иначе алерт-усталость)."""
|
||||
def test_sber_old_period_with_healthy_pull_stays_silent(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Прод 2026-08-12: период старый (72 суток), но загрузка в такте — событий нет.
|
||||
|
||||
Это ровно то состояние, в котором main двенадцатые сутки подряд писал ERROR:
|
||||
возраст там был лагом ПУБЛИКАЦИИ Сбера, а не нашим отставанием. Алерт-усталость
|
||||
от таких событий и делает настоящий отказ незаметным.
|
||||
"""
|
||||
_patch_runs(monkeypatch, sber_mon)
|
||||
db = _FakeMonitorDB(date(2026, 6, 1))
|
||||
db = _FakeMonitorDB(date(2026, 6, 1)) # last_pull = 2026-08-06 (полный прогон)
|
||||
with glitchtip_events() as events:
|
||||
out = sber_mon.check_sber_freshness(
|
||||
db, # type: ignore[arg-type]
|
||||
run_id=2,
|
||||
params={},
|
||||
now=datetime(2026, 6, 20, tzinfo=UTC),
|
||||
now=datetime(2026, 8, 12, 19, 6, tzinfo=UTC),
|
||||
)
|
||||
assert out["alert"] == 0
|
||||
assert event_texts(events) == []
|
||||
|
|
@ -165,7 +191,7 @@ def test_sber_empty_index_becomes_event(monkeypatch: pytest.MonkeyPatch) -> None
|
|||
params={},
|
||||
now=datetime(2026, 8, 6, tzinfo=UTC),
|
||||
)
|
||||
assert any("sber_price_index пуст" in t for t in event_texts(events))
|
||||
assert any("у оценщика нет серии" in t for t in event_texts(events))
|
||||
|
||||
|
||||
def test_deals_empty_becomes_event(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
|
|
|
|||
|
|
@ -331,18 +331,24 @@ def test_fix4_premium_comp_survives_post_weight_clip() -> None:
|
|||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Fix 5a — sber staleness warning
|
||||
# Fix 5a — per-estimate sber staleness warning УДАЛЁН (#2846)
|
||||
#
|
||||
# Guard сравнивал возраст latest периода с sber_index_max_age_days=35. Такой
|
||||
# возраст недостижим по построению (period_month — метка первого числа + лаг
|
||||
# публикации ⇒ пол 46 суток), поэтому на проде warning писался у КАЖДОЙ оценки
|
||||
# и не нёс ни бита. Пара прежних тестов зеленела только на фикстуре с «свежим»
|
||||
# месяцем, которого в реальной серии не бывает. Свежесть теперь мерит одно место —
|
||||
# tasks/sber_freshness_monitor, по отставанию ЗАГРУЗКИ.
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def test_fix5a_stale_sber_logs_warning(caplog: pytest.LogCaptureFixture) -> None:
|
||||
"""_load_sber_index_series логирует warning при stale серии."""
|
||||
def test_fix5a_no_per_estimate_staleness_warning(caplog: pytest.LogCaptureFixture) -> None:
|
||||
"""Серия отдаётся как есть; календарного warning'а в горячем пути больше нет."""
|
||||
import logging
|
||||
|
||||
from app.services.estimator import _load_sber_index_series
|
||||
|
||||
# Серия с единственным месяцем 2 года назад
|
||||
stale_month = date(2024, 1, 1)
|
||||
stale_month = date(2024, 1, 1) # два года назад — прежний guard тут кричал
|
||||
mock_db = MagicMock()
|
||||
mock_db.execute.return_value.mappings.return_value.all.return_value = [
|
||||
{"period_month": stale_month, "index_value_rub_m2": 100_000.0}
|
||||
|
|
@ -351,33 +357,9 @@ def test_fix5a_stale_sber_logs_warning(caplog: pytest.LogCaptureFixture) -> None
|
|||
with caplog.at_level(logging.WARNING, logger="app.services.estimator"):
|
||||
series = _load_sber_index_series(mock_db, region="Свердловская область")
|
||||
|
||||
assert len(series) == 1
|
||||
assert stale_month in series
|
||||
# Warning о staleness должен быть залогирован
|
||||
assert series == {stale_month: 100_000.0}
|
||||
stale_msgs = [r for r in caplog.records if "stale" in r.message.lower()]
|
||||
assert stale_msgs, f"Ожидали warning о stale sber, caplog: {caplog.text}"
|
||||
|
||||
|
||||
def test_fix5a_fresh_sber_no_warning(caplog: pytest.LogCaptureFixture) -> None:
|
||||
"""_load_sber_index_series НЕ логирует warning при свежей серии."""
|
||||
import logging
|
||||
|
||||
from app.services.estimator import _load_sber_index_series
|
||||
|
||||
# Текущий месяц (age=0..30 дней — точно свежее 35-дневного порога).
|
||||
today = datetime.now(tz=UTC).date()
|
||||
fresh_month = today.replace(day=1) # 1-е число ТЕКУЩЕГО месяца
|
||||
mock_db = MagicMock()
|
||||
mock_db.execute.return_value.mappings.return_value.all.return_value = [
|
||||
{"period_month": fresh_month, "index_value_rub_m2": 128_000.0}
|
||||
]
|
||||
|
||||
with caplog.at_level(logging.WARNING, logger="app.services.estimator"):
|
||||
series = _load_sber_index_series(mock_db, region="Свердловская область")
|
||||
|
||||
assert len(series) == 1
|
||||
stale_msgs = [r for r in caplog.records if "stale" in r.message.lower()]
|
||||
assert not stale_msgs, f"Не ожидали stale warning для свежей серии, caplog: {caplog.text}"
|
||||
assert not stale_msgs, f"per-estimate guard вернулся, caplog: {caplog.text}"
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
|
|
|
|||
|
|
@ -1,13 +1,20 @@
|
|||
"""Freshness-монитор данных СберИндекса по max(period_month) — audit п.1.
|
||||
"""Монитор СберИндекса: тревога про ОТСТАВАНИЕ ЗАГРУЗКИ, а не про календарь (#2846).
|
||||
|
||||
Покрывает:
|
||||
1. Чистую логику evaluate_sber_freshness (frozen now, без БД):
|
||||
- fresh: age <= max_age_days (алерта нет);
|
||||
- stale: age > max_age_days (алерт);
|
||||
- граница порога (== max_age_days → нет алерта; +1 день → алерт).
|
||||
2. check_sber_freshness с FakeDB (fresh / stale / empty→mark_failed / кастомный lag).
|
||||
3. Свойства миграции 180 (по образцу test_deals_freshness_monitor).
|
||||
4. Регистрацию в kit product_handlers (registry остаётся зелёным).
|
||||
1. Чистую логику evaluate_sber_freshness (frozen now, без БД): загрузка в такте /
|
||||
загрузка встала / полных прогонов не было вовсе / граница порога.
|
||||
2. check_sber_freshness с FakeDB — прод-реплей 2026-08-12 и двусторонность:
|
||||
при ОДНОМ И ТОМ ЖЕ возрасте периода вердикт меняется вслед за загрузкой.
|
||||
3. Выбор табло тем же порядком, что у оценщика (max() по таблице маскировал бы
|
||||
отставшее табло).
|
||||
4. Свойства миграций 180/212 + регистрацию в kit product_handlers.
|
||||
|
||||
Прод-числа (read-only, 2026-08-12, scrape_runs/scrape_schedules/sber_price_index):
|
||||
latest период табло оценщика real_estate_deals = 2026-06-01 (возраст 72 суток),
|
||||
dinamika-tsen-obyavlenii = 2026-05-01 (103);
|
||||
последний ПОЛНЫЙ прогон загрузки = 2026-08-06 (errors=0, upserted=639);
|
||||
такт загрузки interval_days = 7;
|
||||
прогон id=37 (05-31) — status='done' при {errors: 9, upserted: 0}, за успех НЕ считается.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
|
@ -22,7 +29,6 @@ import pytest
|
|||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||
|
||||
from app.core.config import settings
|
||||
from app.services.product_handlers import build_product_handlers
|
||||
from app.tasks import sber_freshness_monitor as mon
|
||||
|
||||
|
|
@ -30,69 +36,152 @@ _SQL_DIR = Path(__file__).resolve().parents[1] / "data" / "sql"
|
|||
_MIGRATION_180 = _SQL_DIR / "180_seed_sber_freshness_monitor.sql"
|
||||
_MIGRATION_212 = _SQL_DIR / "212_sber_index_pull_weekly.sql"
|
||||
|
||||
# max(period_month) вторичного сегмента = 2026-05-01 (проверено на проде 2026-07-12).
|
||||
_MAY_2026 = date(2026, 5, 1)
|
||||
|
||||
|
||||
# ── evaluate_sber_freshness (чистая логика, frozen now) ───────────────────────
|
||||
# Прод-состояние 2026-08-12.
|
||||
_JUN_2026 = date(2026, 6, 1) # latest табло real_estate_deals — возраст 72 суток
|
||||
_MAY_2026 = date(2026, 5, 1) # latest табло dinamika-tsen-obyavlenii — возраст 103
|
||||
_LAST_FULL_PULL = datetime(2026, 8, 6, 5, 0, tzinfo=UTC) # errors=0, upserted=639
|
||||
_PULL_INTERVAL = 7 # scrape_schedules.default_params.interval_days
|
||||
_PROD_NOW = datetime(2026, 8, 12, 19, 6, tzinfo=UTC) # момент прод-замера
|
||||
|
||||
|
||||
def _now(y: int, m: int, d: int) -> datetime:
|
||||
return datetime(y, m, d, tzinfo=UTC)
|
||||
|
||||
|
||||
def test_fresh_within_max_age() -> None:
|
||||
"""age=30 ≤ max_age_days=60 — свежий, алерта нет."""
|
||||
v = mon.evaluate_sber_freshness(_MAY_2026, _now(2026, 5, 31), max_age_days=60)
|
||||
assert v.stale is False
|
||||
assert v.age_days == 30
|
||||
assert v.latest_period == _MAY_2026
|
||||
# ── evaluate_sber_freshness (чистая логика, frozen now) ───────────────────────
|
||||
|
||||
|
||||
def test_stale_beyond_max_age() -> None:
|
||||
"""Прод-состояние: 2026-05-01 @ 2026-07-12 — age=72 > 60 → алерт."""
|
||||
v = mon.evaluate_sber_freshness(_MAY_2026, _now(2026, 7, 12), max_age_days=60)
|
||||
assert v.stale is True
|
||||
def test_loader_in_cadence_no_alert_even_at_age_72() -> None:
|
||||
"""Прод 2026-08-12: возраст 72, но полный прогон 6 суток назад → молчим.
|
||||
|
||||
После полного прогона наш max(period_month) равен максимуму источника ПО
|
||||
ПОСТРОЕНИЮ (загрузчик тянет всю серию), значит 72 суток — лаг ПУБЛИКАЦИИ Сбера,
|
||||
а не наше отставание.
|
||||
"""
|
||||
v = mon.evaluate_sber_freshness(
|
||||
_JUN_2026,
|
||||
_PROD_NOW,
|
||||
last_complete_pull_at=_LAST_FULL_PULL,
|
||||
pull_interval_days=_PULL_INTERVAL,
|
||||
)
|
||||
assert v.age_days == 72
|
||||
assert v.pull_lag_days == 6
|
||||
assert v.stale is False
|
||||
|
||||
|
||||
def test_loader_stalled_alerts_at_the_same_age() -> None:
|
||||
"""Тот же возраст периода, но полный прогон 20 суток назад → тревога.
|
||||
|
||||
20 суток — реальный разрыв прод-истории (07-17 → 08-06) при пороге 2×7=14.
|
||||
"""
|
||||
v = mon.evaluate_sber_freshness(
|
||||
_JUN_2026,
|
||||
_PROD_NOW,
|
||||
last_complete_pull_at=_now(2026, 7, 23),
|
||||
pull_interval_days=_PULL_INTERVAL,
|
||||
)
|
||||
assert v.age_days == 72 # возраст ТОТ ЖЕ, что в тесте выше
|
||||
assert v.pull_lag_days == 20
|
||||
assert v.max_pull_lag_days == 14
|
||||
assert v.stale is True
|
||||
|
||||
|
||||
def test_no_complete_pull_ever_alerts() -> None:
|
||||
"""Загрузчик умер совсем / не отработал ни разу успешно → тревога, не тишина."""
|
||||
v = mon.evaluate_sber_freshness(
|
||||
_JUN_2026,
|
||||
_PROD_NOW,
|
||||
last_complete_pull_at=None,
|
||||
pull_interval_days=_PULL_INTERVAL,
|
||||
)
|
||||
assert v.stale is True
|
||||
assert v.pull_lag_days == -1
|
||||
|
||||
|
||||
def test_threshold_boundary_exact_no_alert() -> None:
|
||||
"""Ровно на пороге (age == max_age_days) алерта ещё нет (строгое >)."""
|
||||
v = mon.evaluate_sber_freshness(_MAY_2026, _now(2026, 6, 30), max_age_days=60)
|
||||
assert v.age_days == 60
|
||||
"""Ровно на пороге (2 такта) алерта ещё нет — строгое >."""
|
||||
v = mon.evaluate_sber_freshness(
|
||||
_JUN_2026,
|
||||
_LAST_FULL_PULL + timedelta(days=mon.MISSED_PULL_CYCLES * _PULL_INTERVAL),
|
||||
last_complete_pull_at=_LAST_FULL_PULL,
|
||||
pull_interval_days=_PULL_INTERVAL,
|
||||
)
|
||||
assert v.pull_lag_days == 14
|
||||
assert v.stale is False
|
||||
|
||||
|
||||
def test_threshold_boundary_next_day_alert() -> None:
|
||||
"""Порог + 1 день (age=61) — первый алерт."""
|
||||
v = mon.evaluate_sber_freshness(_MAY_2026, _now(2026, 7, 1), max_age_days=60)
|
||||
assert v.age_days == 61
|
||||
def test_threshold_boundary_next_day_alerts() -> None:
|
||||
"""Порог + 1 сутки — первый алерт (два такта подряд пропущены)."""
|
||||
v = mon.evaluate_sber_freshness(
|
||||
_JUN_2026,
|
||||
_LAST_FULL_PULL + timedelta(days=mon.MISSED_PULL_CYCLES * _PULL_INTERVAL + 1),
|
||||
last_complete_pull_at=_LAST_FULL_PULL,
|
||||
pull_interval_days=_PULL_INTERVAL,
|
||||
)
|
||||
assert v.pull_lag_days == 15
|
||||
assert v.stale is True
|
||||
|
||||
|
||||
def test_threshold_follows_pull_cadence() -> None:
|
||||
"""Порог — производная такта загрузки, а не константа: такт 28 → порог 56."""
|
||||
v = mon.evaluate_sber_freshness(
|
||||
_JUN_2026,
|
||||
_PROD_NOW,
|
||||
last_complete_pull_at=_now(2026, 7, 23), # 20 суток
|
||||
pull_interval_days=28,
|
||||
)
|
||||
assert v.max_pull_lag_days == 56
|
||||
assert v.stale is False # при месячном такте 20 суток — норма
|
||||
|
||||
|
||||
# ── check_sber_freshness (FakeDB) ─────────────────────────────────────────────
|
||||
|
||||
|
||||
class _Row:
|
||||
def __init__(self, latest: date | None) -> None:
|
||||
self.latest = latest
|
||||
def __init__(self, **kw: Any) -> None:
|
||||
self.__dict__.update(kw)
|
||||
|
||||
|
||||
class _FakeResult:
|
||||
def __init__(self, latest: date | None) -> None:
|
||||
self._latest = latest
|
||||
def __init__(self, row: _Row | None) -> None:
|
||||
self._row = row
|
||||
|
||||
def first(self) -> _Row:
|
||||
return _Row(self._latest)
|
||||
def first(self) -> _Row | None:
|
||||
return self._row
|
||||
|
||||
|
||||
class _FakeDB:
|
||||
def __init__(self, latest: date | None) -> None:
|
||||
self._latest = latest
|
||||
"""Отвечает на три запроса монитора; пригоден и для старой версии кода.
|
||||
|
||||
Старый монитор спрашивал max(period_month) БЕЗ dashboard-фильтра — на такой
|
||||
запрос отдаём максимум по всем табло, ровно как это делал бы Postgres.
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
latest_by_dash: dict[str, date],
|
||||
last_pull: datetime | None = _LAST_FULL_PULL,
|
||||
interval_days: str | None = str(_PULL_INTERVAL),
|
||||
) -> None:
|
||||
self._latest_by_dash = latest_by_dash
|
||||
self._last_pull = last_pull
|
||||
self._interval_days = interval_days
|
||||
self.rolled_back = False
|
||||
self.asked_dashboards: list[str] = []
|
||||
|
||||
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
|
||||
return _FakeResult(self._latest)
|
||||
sql = str(stmt)
|
||||
params = params or {}
|
||||
if "scrape_runs" in sql:
|
||||
return _FakeResult(_Row(last_pull=self._last_pull))
|
||||
if "scrape_schedules" in sql:
|
||||
return _FakeResult(_Row(interval_days=self._interval_days))
|
||||
if "dash" in params:
|
||||
self.asked_dashboards.append(params["dash"])
|
||||
return _FakeResult(_Row(latest=self._latest_by_dash.get(params["dash"])))
|
||||
# Старый монитор: max(period_month) по всей таблице.
|
||||
latest = max(self._latest_by_dash.values()) if self._latest_by_dash else None
|
||||
return _FakeResult(_Row(latest=latest))
|
||||
|
||||
def rollback(self) -> None:
|
||||
self.rolled_back = True
|
||||
|
|
@ -118,60 +207,138 @@ def _patch_runs(monkeypatch: pytest.MonkeyPatch) -> dict[str, Any]:
|
|||
return calls
|
||||
|
||||
|
||||
def test_check_fresh_marks_done(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
calls = _patch_runs(monkeypatch)
|
||||
db = _FakeDB(_MAY_2026)
|
||||
# @2026-05-31: age=30 ≤ 35+25=60 → нет алерта.
|
||||
out = mon.check_sber_freshness(db, run_id=1, params={}, now=_now(2026, 5, 31)) # type: ignore[arg-type]
|
||||
assert out == {"latest_year": 2026, "latest_month": 5, "age_days": 30, "alert": 0}
|
||||
assert calls["done"] == out
|
||||
assert calls["failed"] is None
|
||||
def _prod_db(last_pull: datetime | None = _LAST_FULL_PULL) -> _FakeDB:
|
||||
"""Прод-состояние 2026-08-12 (оба табло вторички)."""
|
||||
return _FakeDB(
|
||||
{"real_estate_deals": _JUN_2026, "dinamika-tsen-obyavlenii": _MAY_2026},
|
||||
last_pull=last_pull,
|
||||
)
|
||||
|
||||
|
||||
def test_check_stale_marks_done_with_alert(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
def test_prod_replay_healthy_loader_silent_source_no_alert(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
"""КРАСНЫЙ НА main. Прод 2026-08-12: загрузка исправна, источник молчит → тишина.
|
||||
|
||||
На main монитор мерил календарь (72 > 60) и писал ERROR — двенадцатые сутки
|
||||
подряд, при полном прогоне загрузки 08-06. Тревога описывала лаг публикации
|
||||
Сбера, а не наш дефект, и на настоящий отказ загрузчика выглядела бы так же.
|
||||
"""
|
||||
calls = _patch_runs(monkeypatch)
|
||||
db = _FakeDB(_MAY_2026)
|
||||
# @2026-07-12 (прод): age=72 > 60 → алерт.
|
||||
out = mon.check_sber_freshness(db, run_id=2, params={}, now=_now(2026, 7, 12)) # type: ignore[arg-type]
|
||||
assert out["alert"] == 1
|
||||
db = _prod_db()
|
||||
out = mon.check_sber_freshness(db, run_id=1, params={}, now=_PROD_NOW) # type: ignore[arg-type]
|
||||
assert out["age_days"] == 72
|
||||
assert out["latest_month"] == 5
|
||||
# Монитор НЕ падает при алерте — прогон done, а не failed.
|
||||
assert out["alert"] == 0
|
||||
assert calls["done"] == out
|
||||
assert calls["failed"] is None
|
||||
|
||||
|
||||
def test_check_empty_index_marks_failed(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
def test_alert_tracks_loader_not_calendar(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Двусторонность: возраст периода одинаков, вердикт идёт за загрузкой.
|
||||
|
||||
На main оба состояния дают alert=1 (вердикт зависит только от календаря) —
|
||||
сторож не умеет зеленеть, что и было исходным дефектом.
|
||||
"""
|
||||
_patch_runs(monkeypatch)
|
||||
healthy = mon.check_sber_freshness(
|
||||
_prod_db(last_pull=_LAST_FULL_PULL), # type: ignore[arg-type]
|
||||
run_id=2,
|
||||
params={},
|
||||
now=_PROD_NOW,
|
||||
)
|
||||
stalled = mon.check_sber_freshness(
|
||||
_prod_db(last_pull=_now(2026, 7, 23)), # 20 суток назад > 14 # type: ignore[arg-type]
|
||||
run_id=3,
|
||||
params={},
|
||||
now=_PROD_NOW,
|
||||
)
|
||||
assert healthy["age_days"] == stalled["age_days"] == 72
|
||||
assert (healthy["alert"], stalled["alert"]) == (0, 1)
|
||||
|
||||
|
||||
def test_dead_loader_never_pulled_marks_alert(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Загрузчик умер совсем (ни одного полного прогона) → монитор НЕ молчит.
|
||||
|
||||
Прогон id=37 со status='done' при {errors: 9, upserted: 0} за успех не идёт —
|
||||
SQL требует errors=0 AND upserted>0, поэтому «полных прогонов не было» здесь
|
||||
ровно то состояние, что дал бы прод с одним лишь id=37.
|
||||
"""
|
||||
_patch_runs(monkeypatch)
|
||||
out = mon.check_sber_freshness(
|
||||
_prod_db(last_pull=None), # type: ignore[arg-type]
|
||||
run_id=4,
|
||||
params={},
|
||||
now=_PROD_NOW,
|
||||
)
|
||||
assert out["alert"] == 1
|
||||
assert out["pull_lag_days"] == -1
|
||||
|
||||
|
||||
def test_asks_estimator_dashboard_first(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Табло — то же и в том же порядке, что берёт оценщик (не max() по таблице)."""
|
||||
_patch_runs(monkeypatch)
|
||||
db = _prod_db()
|
||||
out = mon.check_sber_freshness(db, run_id=5, params={}, now=_PROD_NOW) # type: ignore[arg-type]
|
||||
assert db.asked_dashboards[0] == "real_estate_deals"
|
||||
assert (out["latest_year"], out["latest_month"]) == (2026, 6)
|
||||
|
||||
|
||||
def test_falls_back_to_next_dashboard_when_first_empty(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
"""Первое табло пусто → берём следующее, как и оценщик."""
|
||||
_patch_runs(monkeypatch)
|
||||
db = _FakeDB({"dinamika-tsen-obyavlenii": _MAY_2026})
|
||||
out = mon.check_sber_freshness(db, run_id=6, params={}, now=_PROD_NOW) # type: ignore[arg-type]
|
||||
assert out["latest_month"] == 5
|
||||
assert out["age_days"] == 103
|
||||
|
||||
|
||||
def test_empty_index_marks_failed(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Ни одного табло с данными — оценивать нечего, это сбой монитора."""
|
||||
calls = _patch_runs(monkeypatch)
|
||||
db = _FakeDB(None)
|
||||
out = mon.check_sber_freshness(db, run_id=3, params={}, now=_now(2026, 7, 12)) # type: ignore[arg-type]
|
||||
out = mon.check_sber_freshness(_FakeDB({}), run_id=7, params={}, now=_PROD_NOW) # type: ignore[arg-type]
|
||||
assert out["alert"] == 0
|
||||
assert calls["done"] is None
|
||||
assert calls["failed"] is not None
|
||||
|
||||
|
||||
def test_check_reads_lag_from_params(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
def test_legacy_lag_allowance_param_is_ignored(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Мёртвая ручка default_params.lag_allowance_days не может вернуть календарь.
|
||||
|
||||
Строка монитора в проде всё ещё несёт {"lag_allowance_days": 25} (миграция 180).
|
||||
С любым её значением вердикт один и тот же — порог берётся из такта загрузки.
|
||||
"""
|
||||
_patch_runs(monkeypatch)
|
||||
db = _FakeDB(_MAY_2026)
|
||||
# lag=0 → порог = sber_index_max_age_days (35) → @2026-07-12 (age=72) просрочено.
|
||||
out = mon.check_sber_freshness(
|
||||
db, run_id=4, params={"lag_allowance_days": 0}, now=_now(2026, 7, 12)
|
||||
) # type: ignore[arg-type]
|
||||
assert out["alert"] == 1
|
||||
base = mon.check_sber_freshness(_prod_db(), run_id=8, params={}, now=_PROD_NOW) # type: ignore[arg-type]
|
||||
tweaked = mon.check_sber_freshness(
|
||||
_prod_db(), # type: ignore[arg-type]
|
||||
run_id=9,
|
||||
params={"lag_allowance_days": 0},
|
||||
now=_PROD_NOW,
|
||||
)
|
||||
assert base["alert"] == tweaked["alert"] == 0
|
||||
|
||||
|
||||
def test_check_default_threshold_uses_setting_plus_lag(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Дефолтный порог = sber_index_max_age_days + DEFAULT_LAG_ALLOWANCE_DAYS."""
|
||||
def test_threshold_read_from_pull_schedule(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Порог читается из строки загрузчика: такт 28 → порог 56, тревоги нет."""
|
||||
_patch_runs(monkeypatch)
|
||||
db = _FakeDB(_MAY_2026)
|
||||
threshold = settings.sber_index_max_age_days + mon.DEFAULT_LAG_ALLOWANCE_DAYS
|
||||
# Ровно на пороге (age == threshold) — алерта нет; +1 день — алерт.
|
||||
exact = datetime(2026, 5, 1, tzinfo=UTC) + timedelta(days=threshold)
|
||||
out_exact = mon.check_sber_freshness(db, run_id=5, params={}, now=exact) # type: ignore[arg-type]
|
||||
assert out_exact["age_days"] == threshold
|
||||
assert out_exact["alert"] == 0
|
||||
out_over = mon.check_sber_freshness(db, run_id=6, params={}, now=exact + timedelta(days=1)) # type: ignore[arg-type]
|
||||
assert out_over["alert"] == 1
|
||||
db = _FakeDB(
|
||||
{"real_estate_deals": _JUN_2026},
|
||||
last_pull=_now(2026, 7, 23), # 20 суток
|
||||
interval_days="28",
|
||||
)
|
||||
out = mon.check_sber_freshness(db, run_id=10, params={}, now=_PROD_NOW) # type: ignore[arg-type]
|
||||
assert out["max_pull_lag_days"] == 56
|
||||
assert out["alert"] == 0
|
||||
|
||||
|
||||
def test_missing_schedule_row_falls_back_to_default(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Строки/ключа нет — берём DEFAULT_PULL_INTERVAL_DAYS, а не падаем."""
|
||||
_patch_runs(monkeypatch)
|
||||
db = _FakeDB({"real_estate_deals": _JUN_2026}, interval_days=None)
|
||||
out = mon.check_sber_freshness(db, run_id=11, params={}, now=_PROD_NOW) # type: ignore[arg-type]
|
||||
assert out["max_pull_lag_days"] == mon.MISSED_PULL_CYCLES * mon.DEFAULT_PULL_INTERVAL_DAYS
|
||||
|
||||
|
||||
# ── Миграция 180 ──────────────────────────────────────────────────────────────
|
||||
|
|
@ -208,22 +375,18 @@ def test_migration_180_window_9_to_10_utc() -> None:
|
|||
assert re.search(r"\b10\b", sql), "window_end_hour 10 missing"
|
||||
|
||||
|
||||
def test_migration_180_lag_allowance_25() -> None:
|
||||
sql = _MIGRATION_180.read_text("utf-8")
|
||||
assert "lag_allowance_days" in sql
|
||||
assert "25" in sql
|
||||
|
||||
|
||||
def test_migration_180_no_psycopg_trap() -> None:
|
||||
sql = _MIGRATION_180.read_text("utf-8")
|
||||
assert not re.search(r":\w+::", sql)
|
||||
|
||||
|
||||
# ── Миграция 212: такт загрузки не должен пересекать порог монитора ───────────
|
||||
# ── Миграция 212: такт загрузки = источник порога ─────────────────────────────
|
||||
#
|
||||
# Прод-разбор (ревью PR #2681): загрузка раз в 28 дней давала возраст-пилу 46..74
|
||||
# при пороге 60 — тревога срабатывала 14 суток из 28 БЕЗ всякого застоя источника.
|
||||
# Тест держит инвариант: потолок возраста (пол + такт загрузки) < порога монитора.
|
||||
# Прежний инвариант («пол возраста + такт < календарного порога монитора») снят
|
||||
# вместе с календарным порогом: прод его ОПРОВЕРГ — 2026-08-12 возраст 72 при
|
||||
# полном прогоне шестидневной давности, потолок 53 держался бы только если бы
|
||||
# источник публиковал строго помесячно. Остаётся то, что проверяемо: фолбэк кода
|
||||
# не должен расходиться с тактом, который сеет миграция.
|
||||
|
||||
|
||||
def test_migration_212_makes_pull_cadence_weekly() -> None:
|
||||
|
|
@ -234,23 +397,22 @@ def test_migration_212_makes_pull_cadence_weekly() -> None:
|
|||
assert not re.search(r":\w+::", sql) # psycopg v3: только CAST(:x AS type)
|
||||
|
||||
|
||||
def test_pull_cadence_leaves_margin_under_monitor_threshold() -> None:
|
||||
"""Инвариант: пол возраста + такт загрузки < порога монитора.
|
||||
|
||||
Пол = 46 суток (прод 2026-07-17: загрузка принесла 2026-06-01). Порог =
|
||||
sber_index_max_age_days + lag_allowance. При такте 7: 46+7=53 < 60 — запас
|
||||
7 суток. При прежних 28: 46+28=74 > 60 — тревога каждый цикл, что и наблюдали.
|
||||
"""
|
||||
def test_default_pull_interval_matches_migration_212() -> None:
|
||||
"""Фолбэк монитора == такт из миграции, иначе порог тихо разъедется с загрузкой."""
|
||||
interval_days = int(
|
||||
re.search(r'"interval_days":\s*(\d+)', _MIGRATION_212.read_text("utf-8")).group(1)
|
||||
)
|
||||
observed_floor_days = 46
|
||||
threshold = settings.sber_index_max_age_days + mon.DEFAULT_LAG_ALLOWANCE_DAYS
|
||||
assert observed_floor_days + interval_days < threshold, (
|
||||
f"такт {interval_days}д даёт потолок возраста "
|
||||
f"{observed_floor_days + interval_days}д при пороге {threshold}д — "
|
||||
"монитор снова будет мерить наш такт, а не застой источника"
|
||||
re.search(r'"interval_days":\s*(\d+)', _MIGRATION_212.read_text("utf-8")).group(1) # type: ignore[union-attr]
|
||||
)
|
||||
assert mon.DEFAULT_PULL_INTERVAL_DAYS == interval_days
|
||||
|
||||
|
||||
# ── Один порог, а не два ──────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_no_second_calendar_threshold_in_settings() -> None:
|
||||
"""#2846: sber_index_max_age_days удалён — второму порогу неоткуда взяться."""
|
||||
from app.core.config import settings
|
||||
|
||||
assert not hasattr(settings, "sber_index_max_age_days")
|
||||
|
||||
|
||||
# ── Регистрация в kit registry ─────────────────────────────────────────────────
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue