"""Монитор ОТСТАВАНИЯ ЗАГРУЗКИ СберИндекса (не календарного возраста периода). ЧТО БЫЛО НЕ ТАК (замер на проде 2026-08-12, #2846). Монитор мерил `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% времени — ноль бит. Порог 60 у монитора не лучше: он лежит ВНУТРИ рабочего диапазона, поэтому монитор мерил не источник, а нашу же пилу. Миграция 212 (такт 28 → 7) обещала потолок возраста ≈46+7=53 < 60. Прод это ОПРОВЕРГ: 2026-08-12 возраст 72 при ПОЛНОМ прогоне загрузки шестидневной давности (08-06, errors=0, upserted=639) — источник просто не опубликовал июль. Двенадцать суток подряд (08-01 … 08-12) монитор писал ERROR при исправной загрузке. Потолок 53 держится, только если источник публикует строго помесячно; он не публикует. ЧТО МЕРИМ ТЕПЕРЬ. Загрузчик тянет ВСЮ серию (limit=1000&offset=0, отсечки по периоду нет), поэтому после прогона с errors=0 AND upserted>0 наш max(period_month) РАВЕН максимуму источника ПО ПОСТРОЕНИЮ. Значит вопрос «отстали ли мы» — это вопрос «давно ли был последний ЗАВЕДОМО ПОЛНЫЙ прогон», и он не зависит от возраста периода: последний полный прогон свежий → наш max == max источника → источник не публиковал, молчание ПРАВИЛЬНОЕ (возраст = лаг источника); последний полный прогон старый → мы не забрали → тревога про ЗАГРУЗЧИК. ЛОВУШКА: `status='done'` НЕ означает успех — прогон id=37 (2026-05-31) имеет {errors: 9, upserted: 0} и статус done. Успех = errors=0 AND upserted>0 (все 9 серий 3 табло × 3 региона прошли: errors — счётчик по всему прогону). ПОРОГ — не круглое число, а такт самой загрузки: `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 import logging from dataclasses import dataclass from datetime import UTC, date, datetime from sqlalchemy import text from sqlalchemy.orm import Session 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_PULL_INTERVAL_DAYS", "MISSED_PULL_CYCLES", "SBER_FRESHNESS_PULL_SOURCE", "SberFreshnessVerdict", "check_sber_freshness", "evaluate_sber_freshness", ] # Джоба-загрузчик, чей такт и успешность мы и мониторим. SBER_FRESHNESS_PULL_SOURCE = "sber_index_pull" # Сколько тактов загрузки подряд можно пропустить до тревоги. 1 = разовый сбой/сдвиг # окна (поглощаем), 2 = загрузка встала (алерт). MISSED_PULL_CYCLES = 2 # Фолбэк, если в 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: """Вердикт: отстала ли ЗАГРУЗКА СберИндекса от собственного такта.""" latest_period: date 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, *, last_complete_pull_at: datetime | None, pull_interval_days: int, ) -> SberFreshnessVerdict: """Чистая логика: отстала ли загрузка от собственного такта. stale = полных прогонов не было ВООБЩЕ, либо последний старше MISSED_PULL_CYCLES × pull_interval_days. Возраст периода считается и кладётся в вердикт как НАБЛЮДЕНИЕ, но на вердикт не влияет: после полного прогона наш max(period_month) равен максимуму источника по построению, и его возраст — это лаг ПУБЛИКАЦИИ, на который мы повлиять не можем. """ age_days = (now.date() - latest_period).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, 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]: """Проверить, не отстала ли загрузка СберИндекса, и алертить при отставании. Sync (вызывается scheduler-триггером в executor, как check_deals_freshness). Читает: latest период табло оценщика, время последнего ПОЛНОГО прогона sber_index_pull, такт загрузки из scrape_schedules. Вердикт — чистой функцией. `params` больше ничего не настраивает: порог берётся из такта самой загрузки (унаследованный default_params.lag_allowance_days=25 монитора игнорируется — он кодировал мёртвый календарный порог). `now` инъектируется в тестах. Returns counters {latest_year, latest_month, age_days, pull_lag_days, max_pull_lag_days, alert}. mark_failed только если у оценщика нет серии вообще (нечего оценивать); при алерте прогон помечается done (это монитор, не сбой джобы). """ 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) found = _load_estimator_dashboard(db) if found is None: # ERROR (#2674): монитор не может выполнить свою работу вовсе — это сбой, # а не наблюдение. mark_failed ниже виден только стрик-алерту (3 подряд), # а монитор ходит раз в сутки — три дня молчания на пустом бенчмарке. logger.error( "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 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). logger.error( "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, dashboard, ) else: logger.info( "sber freshness: загрузка в такте — последний полный прогон %d суток назад " "(≤ порога %d). max(period_month)=%s (табло %s, возраст %d суток) равен " "максимуму источника по построению: возраст = лаг ПУБЛИКАЦИИ источника, " "не наше отставание — алерта нет", verdict.pull_lag_days, verdict.max_pull_lag_days, latest, dashboard, verdict.age_days, ) runs_mod.mark_done(db, run_id, counters) logger.info( "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 except Exception as exc: logger.exception("check_sber_freshness run_id=%d failed", run_id) try: db.rollback() except Exception: pass runs_mod.mark_failed(db, run_id, str(exc)[:1000], counters) raise