"""Продуктовые Handler'ы для kit-scheduler (#2192, strangler-миграция). `build_product_handlers(ctx)` собирает реестр НЕ-sweep source'ов, тело которых осталось в `app` (rosreestr_dkp / sber_index / deactivate_stale_* / *_backfill / cadastral_geo_match / house_* / proxy_healthcheck / …). Kit-native sweep-оркестраторы (avito/yandex/cian/domclick full-load) регистрируются самим `build_registry` через `_default_kit_handlers` — здесь их НЕТ. Исключение: `domclick_city_sweep` (#3264) — kit-native тело осталось прежним (`run_domclick_city_sweep`), но здесь переопределён Handler, который ДО вызова читает БД за куки сессии (kit не имеет права на app.*/БД, см. докстринг `_job_domclick_city_sweep` ниже); `build_registry` явно допускает такое переопределение продуктовым ключом. Дизайн-инвариант: НЕ дублируем логику. Каждый `job` переиспользует существующую боевую джоб-функцию (`app.services.scheduler.import_rosreestr_dkp`, `app.tasks.*`, `app.services.*`) — то же тело, что крутит боевой `trigger_*_run`. Kit `_dispatch` уже делает claim + fresh-session + spawn + close, поэтому `job` содержит ТОЛЬКО работу (без claim/spawn-boilerplate). Run-lifecycle (mark_done/mark_failed/heartbeat) для джоб, что не владеют им сами, делается через инжектированный `ctx.runs`. Ship-dark: модуль импортируется лишь когда `settings.use_kit_scheduler` True (scheduler_main._run_kit_scheduler). Боевой `app.services.scheduler` не тронут. """ from __future__ import annotations import asyncio import logging from datetime import UTC, datetime, timedelta from typing import TYPE_CHECKING, Any from scraper_kit.orchestration import runs as kit_runs from scraper_kit.orchestration.scheduler import ( Handler, reschedule_after_minutes, ) from scraper_kit.orchestration.scheduler import ( _defer_next_run_at as kit_defer_next_run_at, ) from scraper_kit.orchestration.scheduler import ( _pick_resume as kit_pick_resume, ) if TYPE_CHECKING: from scraper_kit.orchestration.scheduler import SchedulerContext from sqlalchemy.orm import Session logger = logging.getLogger(__name__) # ── cian_history_backfill — cookie-gated backfill ──────────────────────────── # Машиночитаемые причины пропуска (#2658) — пишутся в scrape_runs.error строки # со status='skipped'. Отделены от kit-причин (already_running и т.п.): по слагу # видно, встал ли сбор из-за кук или из-за конкурентного прогона. SKIP_CIAN_COOKIES_MISSING = "cian_cookies_missing" SKIP_CIAN_COOKIES_EXPIRED = "cian_cookies_expired" SKIP_CIAN_COOKIES_INVALID = "cian_cookies_invalid" def _alert_cian_cookies(source: str, detail: str) -> None: """Громкий алерт «сбор встал из-за кук» — logger.error, НЕ capture_message(warning). В scraper-контейнере GlitchTip поднят с LoggingIntegration(event_level=ERROR) (scheduler_main.py) — ERROR-запись сама становится событием, а прежний `capture_message(..., level="warning")` до этого уровня не дотягивал (и стоял в недостижимой ветке, см. докстринг _cian_pre_claim). Заодно причина остаётся в docker-логах и в строке scrape_runs, которая переживает редеплой. """ logger.error( "scheduler: %s пропущен — %s. Перезалейте куки Циана через админку " "(до этого backfill истории стоит)", source, detail, ) async def _cian_pre_claim(db: Session, schedule_row: dict[str, Any], ctx: SchedulerContext) -> bool: """Pre-claim gate: проверить наличие/валидность cian-cookies ДО claim (#1522). Cookies отсутствуют/протухли → пишем строку прогона status='skipped' с причиной, двигаем next_run_at на следующее окно и skip (иначе get_due_schedules переотбирает schedule каждые 60с и verify_session долбит Cian круглосуточно). #2658 — что было не так. Первая ветка (load_session вернул None) молчала: warning в docker-лог, сдвиг next_run_at, `return False`. Ни строки в scrape_runs, ни изменения last_run_at — снаружи 37 дней простоя выглядели как «всё по расписанию». Sentry-алерт стоял во ВТОРОЙ ветке (verify_session вернул None), до которой на протухших куках исполнение не доходит НИКОГДА: load_session сам фильтрует expires_at_estimate > NOW() и отдаёт None ещё в первой. Теперь громко в обеих + предупреждение ЗАРАНЕЕ, пока куки ещё валидны (COOKIE_EXPIRY_WARN_DAYS) — обновление кук ручное, ему нужен запас. """ from app.services.cian_session import ( COOKIE_EXPIRY_WARN_DAYS, load_session, session_expires_at, verify_session, ) source: str = schedule_row["source"] now = datetime.now(tz=UTC) cookies = load_session(db) if cookies is None: expires_at = session_expires_at(db) if expires_at is None: reason, detail = SKIP_CIAN_COOKIES_MISSING, "кук Циана нет в БД" elif expires_at <= now: reason = SKIP_CIAN_COOKIES_EXPIRED detail = ( f"куки Циана протухли {expires_at:%Y-%m-%d} ({(now - expires_at).days} дн. назад)" ) else: reason = SKIP_CIAN_COOKIES_INVALID detail = "куки Циана помечены невалидными (last_invalid_at)" _alert_cian_cookies(source, detail) kit_runs.mark_skipped(db, source=source, reason=reason, details=detail) kit_defer_next_run_at(db, schedule_row) return False state = await verify_session(cookies) if state is None: # verify вернул именно None (401 / isAuthenticated=false) — куки числятся # валидными по сроку, но Циан их не принимает. Sentinel-ответы (бан / источник # недоступен / сменилась вёрстка) сюда НЕ попадают, они truthy — см. cian_session. detail = "Циан не принимает куки (разлогин)" _alert_cian_cookies(source, detail) kit_runs.mark_skipped(db, source=source, reason=SKIP_CIAN_COOKIES_INVALID, details=detail) kit_defer_next_run_at(db, schedule_row) return False # Куки рабочие — предупреждаем, пока есть время их обновить без простоя сбора. # valid_only=True: срок ИМЕННО той записи, которую взял load_session (при нескольких # аккаунтах свежайшая-любая может быть чужой протухшей строкой). expires_at = session_expires_at(db, valid_only=True) if expires_at is not None and expires_at - now <= timedelta(days=COOKIE_EXPIRY_WARN_DAYS): logger.error( "scheduler: куки Циана протухнут %s (осталось %.1f дн.) — обновите заранее, " "иначе %s встанет молча", expires_at.date().isoformat(), (expires_at - now).total_seconds() / 86400, source, ) return True async def _job_cian_history_backfill( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.services.scheduler import _execute_cian_backfill await _execute_cian_backfill(db, run_id=run_id, params=params) # ── rosreestr_dkp_import — sync FDW-импорт в executor ───────────────────────── async def _job_rosreestr_dkp( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.services.scheduler import import_rosreestr_dkp loop = asyncio.get_event_loop() await loop.run_in_executor(None, import_rosreestr_dkp, db, run_id, params) # ── listing_source_snapshot — sync DB-snapshot в executor ──────────────────── # params прокинуты (#2607) — snapshot_listing_sources теперь читает budget_sec из # default_params (SET LOCAL statement_timeout, см. app/tasks/listing_source_snapshot.py). async def _job_listing_source_snapshot( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.listing_source_snapshot import snapshot_listing_sources loop = asyncio.get_event_loop() await loop.run_in_executor(None, snapshot_listing_sources, db, run_id, params) # ── asking_to_sold_ratio_refresh — sync re-derive в executor ───────────────── async def _job_asking_to_sold_ratio( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.asking_to_sold_ratio import recompute_asking_to_sold_ratios loop = asyncio.get_event_loop() await loop.run_in_executor(None, recompute_asking_to_sold_ratios, db, run_id) # ── deal_city_price_bands_refresh — sync tier-aware re-derive в executor ────── async def _job_deal_city_price_bands_refresh( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.deal_city_price_bands_refresh import refresh_deal_city_price_bands loop = asyncio.get_event_loop() await loop.run_in_executor(None, refresh_deal_city_price_bands, db, run_id) # ── refresh_search_matview — REFRESH MATVIEW CONCURRENTLY (own connection) ──── async def _job_refresh_search_matview( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.refresh_search_matview import refresh_search_matview loop = asyncio.get_event_loop() try: # refresh_search_matview() держит собственное autocommit-соединение # (REFRESH CONCURRENTLY нельзя внутри транзакции). await loop.run_in_executor(None, refresh_search_matview) ctx.runs.mark_done(db, run_id, {}) except Exception: logger.exception("scheduler: refresh_search_matview crashed run_id=%d", run_id) ctx.runs.mark_failed(db, run_id, "refresh_search_matview failed", {}) # ── yandex_address_backfill — async, owns lifecycle ────────────────────────── async def _job_yandex_address_backfill( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.yandex_address_backfill import run_yandex_address_backfill await run_yandex_address_backfill(db, run_id=run_id, params=params) # ── deactivate_stale_* — sync UPDATE в executor (wildcard-семейство) ────────── async def _job_deactivate_stale( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.core.config import settings as _settings from app.tasks.deactivate_stale_avito import ( CAP_MULT, DEFAULT_FLOOR_DROP_RATIO, DEFAULT_MAX_DEACTIVATED, DEFAULT_MIN_CONFIRMATIONS, DEFAULT_MIN_FLOOR_PAIRS, DEFAULT_REVISIT_FLOOR_QUANTILE, deactivate_stale_listings, ) listing_source: str = params.get("listing_source", "avito") ttl_days: int = params.get("ttl_days", _settings.avito_stale_ttl_days) segments: list[str] | None = params.get("segments") staleness_column: str = params.get("staleness_column", "last_seen_at") # Гейт по здоровью сбора (#2659) включён по умолчанию: незасеянное расписание # получает страховочный порог, а не «деактивируй вслепую». Посчитанные по # источнику пороги приходят из default_params (миграция 219). min_confirmations: int = params.get("min_confirmations", DEFAULT_MIN_CONFIRMATIONS) # Пол TTL по измеренному циклу переобхода (#2659) — тоже включён по умолчанию: # незасеянное расписание не должно снимать объявления по порогу ниже собственного # хвоста обхода. Снять ручку вручную: revisit_floor_quantile = 0. revisit_floor_quantile: float = params.get( "revisit_floor_quantile", DEFAULT_REVISIT_FLOOR_QUANTILE ) # Пустой (NULL) listing_segment -- легаси-строки до миграции 011 + жертвы # отсутствующего COALESCE в ON CONFLICT (base.py upsert никогда не перезаписывает # listing_segment на повторном скрейпе). Отдельный явный предикат IS NULL, а не # элемент :segments (ANY(...) никогда не матчит NULL) -- см. deactivate_stale_avito.py. null_segment_only: bool = params.get("null_segment_only", False) # Потолок эффективного TTL (см. CAP_MULT в deactivate_stale_avito.py) — множитель, # а не голая константа: источник с непропорционально длинным хвостом переобхода # относительно своего ttl_days переопределяет его через default_params (ключ # "cap_mult"), не трогая дефолт для остальных источников. cap_mult: float = params.get("cap_mult", CAP_MULT) # Гейт деградации пола (PR-B, #2659 продолжение) -- включён по умолчанию, тот # же принцип, что у min_confirmations/revisit_floor_quantile выше: незасеянное # расписание получает страховку, а не «деактивируй вслепую». min_floor_pairs: int = params.get("min_floor_pairs", DEFAULT_MIN_FLOOR_PAIRS) floor_drop_ratio: float = params.get("floor_drop_ratio", DEFAULT_FLOOR_DROP_RATIO) # Аварийный (не рабочий) потолок объёма снятия за один прогон -- см. # DEFAULT_MAX_DEACTIVATED в deactivate_stale_avito.py. max_deactivated: int = params.get("max_deactivated", DEFAULT_MAX_DEACTIVATED) loop = asyncio.get_event_loop() await loop.run_in_executor( None, lambda: deactivate_stale_listings( db, run_id, listing_source=listing_source, ttl_days=ttl_days, segments=segments, staleness_column=staleness_column, min_confirmations=min_confirmations, revisit_floor_quantile=revisit_floor_quantile, null_segment_only=null_segment_only, cap_mult=cap_mult, min_floor_pairs=min_floor_pairs, floor_drop_ratio=floor_drop_ratio, max_deactivated=max_deactivated, ), ) # ── sber_index_pull — async, owns lifecycle ────────────────────────────────── async def _job_sber_index_pull( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.sber_index_pull import run_sber_index_pull await run_sber_index_pull(db, run_id=run_id, params=params) # ── rosreestr_quarter_poll — async, owns lifecycle ─────────────────────────── async def _job_rosreestr_quarter_poll( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.rosreestr_quarter_poll import run_rosreestr_quarter_poll await run_rosreestr_quarter_poll(db, run_id=run_id, params=params) # ── deals_freshness_monitor — sync DB-only freshness check в executor ───────── async def _job_deals_freshness_monitor( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.deals_freshness_monitor import check_deals_freshness loop = asyncio.get_event_loop() await loop.run_in_executor(None, check_deals_freshness, db, run_id, params) # ── landing_stats_refresh — sync DB-only пересчёт витрины в executor ───────── async def _job_landing_stats( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.landing_stats import refresh_landing_stats loop = asyncio.get_event_loop() await loop.run_in_executor(None, refresh_landing_stats, db, run_id, params) # ── landing_showcase_deals — sync пересчёт витрины сделок в executor ───────── async def _job_landing_showcase_deals( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: """Пересчёт витрины сделок публичного лэндинга (#3469). ЛАЙФСАЙКЛ ПРОГОНА ВЕДЁТ HANDLER, а не задача. `refresh_landing_showcase_deals` писалась под ручной запуск (`python -m app.tasks.landing_showcase_deals`) и про `run_id` ничего не знает — тот же случай, что у `_job_refresh_search_matview`, и решается так же: done/failed ставим здесь. Параметры берём ИЗ РАСПИСАНИЯ только те, что в нём есть: дефолты живут в сигнатуре задачи, и повтор их здесь дал бы два места, которые разъедутся. """ from app.tasks.landing_showcase_deals import refresh_landing_showcase_deals kwargs = {k: params[k] for k in ("sample", "since", "limit", "city") if k in params} loop = asyncio.get_event_loop() try: counters = await loop.run_in_executor( None, lambda: refresh_landing_showcase_deals(db, **kwargs) ) ctx.runs.mark_done(db, run_id, counters) except Exception: logger.exception("scheduler: landing_showcase_deals crashed run_id=%d", run_id) ctx.runs.mark_failed(db, run_id, "landing_showcase_deals failed", {}) # ── sber_freshness_monitor — sync DB-only freshness check в executor ────────── async def _job_sber_freshness_monitor( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.sber_freshness_monitor import check_sber_freshness loop = asyncio.get_event_loop() await loop.run_in_executor(None, check_sber_freshness, db, run_id, params) # ── newbuilding_enrich — async, owns lifecycle ─────────────────────────────── async def _job_newbuilding_enrich( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.newbuilding_enrich_backfill import run_newbuilding_enrich await run_newbuilding_enrich(db, run_id=run_id, params=params) # ── yandex_newbuilding_sweep — async, lifecycle в job ──────────────────────── async def _job_yandex_newbuilding_sweep( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.yandex_newbuilding_sweep import enrich_yandex_newbuilding_sweep try: result = await enrich_yandex_newbuilding_sweep( db, limit=int(params.get("limit", 5)), request_delay_sec=float(params.get("request_delay_sec", 8.0)), city=str(params.get("city", "ekaterinburg")), ) counters = result.to_dict() # #2860: прогон, который обработал дома и не разрешил НИ ОДНОГО, успешным # называть нельзя. Четырнадцать таких прогонов подряд (16.07-17.08.2026) # стояли `done` с пустым error_text — механизм исполнялся, счётчик был # честный, вывод из него не делал никто, и очередь тихо росла 351 → 397. # # НЕ подгоняем счётчик: если дома действительно не разрешаются, честный # исход — назвать прогон неуспешным, а не дотянуть succeeded до ненуля. processed = int(counters.get("processed") or 0) succeeded = int(counters.get("succeeded") or 0) if result.no_proxy_stop: # #3197: прогон оборван на пустом пуле — к площадке не ходили вовсе. # Это отказ нашей инфраструктуры, а не «ЖК не разрешились»: называть # такой прогон успешным нельзя, и причина должна быть отличима. # # INFO, как у соседей (scheduler.py:175, avito_detail_backfill.py:1050, # domclick:626): причина уже записана в mark_failed + counters.no_proxy_stop=1, # а ERROR ставил её в один разряд с падением задачи (logger.exception ниже). logger.info( "yandex_newbuilding_sweep run_id=%d: пул прокси пуст — прогон оборван " "(обработано %d)", run_id, processed, ) ctx.runs.mark_failed( db, run_id, "пул прокси пуст — прогон оборван, к площадке не ходили", counters, ) elif processed > 0 and succeeded == 0: logger.warning( "yandex_newbuilding_sweep run_id=%d: обработано %d, разрешено 0 — " "помечаю прогон неуспешным", run_id, processed, ) ctx.runs.mark_failed( db, run_id, f"обработано {processed} домов, разрешено 0 — полный отказ разрешения slug", counters, ) else: ctx.runs.mark_done(db, run_id, counters) except Exception: logger.exception("scheduler: enrich_yandex_newbuilding_sweep crashed run_id=%d", run_id) try: ctx.runs.mark_failed(db, run_id, "crashed", {}) except Exception: pass # ── geocode_missing_listings — async, owns lifecycle ───────────────────────── async def _job_geocode_missing_listings( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.geocode_missing import run_geocode_missing_listings await run_geocode_missing_listings(db, run_id=run_id, params=params) # ── avito_detail_backfill — async, owns lifecycle ──────────────────────────── async def _job_avito_detail_backfill( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.avito_detail_backfill import run_avito_detail_backfill await run_avito_detail_backfill(db, run_id=run_id, params=params) # ── yandex_detail_backfill — async, owns lifecycle ─────────────────────────── async def _job_yandex_detail_backfill( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.yandex_detail_backfill import run_yandex_detail_backfill await run_yandex_detail_backfill(db, run_id=run_id, params=params) # ── domclick_detail_backfill — async, owns lifecycle (issue #2000 Layer B) ─── async def _job_domclick_detail_backfill( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.domclick_detail_backfill import run_domclick_detail_backfill await run_domclick_detail_backfill(db, run_id=run_id, params=params) # ── domclick_city_sweep — override kit-native с инъекцией куки сессии (#3264) ─ # Kit-native `_job_domclick_city_sweep` (scraper_kit.orchestration.scheduler) ходит на # bff-search-web.domclick.ru БЕЗ кук: kit не имеет права импортировать app.* / читать БД # (strangler-инвариант #2133). Каждый запрос свипа поэтому был обязан решать QRATOR PoW # с нуля — прод 30.08 (прогон 5330): 9/9 запросов зависли на challenge, 0 листингов. # `domclick_detail_backfill` (добор карточек) уже решает эту задачу инъекцией того же # снимка (domclick_session.load_session) — здесь тот же приём для свипа. # # Регистрируем ЗДЕСЬ (app-side) под тем же ключом "domclick_city_sweep", а не правим # kit: build_registry явно допускает переопределение kit-native продуктовым Handler'ом # (см. докстринг build_registry — "последнее слово за продуктом"), а строгий запрет на # app.* внутри kit при этом не нарушается — БД читает только этот модуль. async def _job_domclick_city_sweep( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from scraper_kit.orchestration.pipeline import run_domclick_city_sweep from app.services import domclick_session as domclick_session_svc # None — сессии в БД нет либо протухла (load_session сам фильтрует # expires_at_estimate > NOW()). Это НЕ авария: свип продолжает работать без # инъекции, ровно как до #3264, только логируем факт один раз для видимости. cookies = domclick_session_svc.load_session(db) if cookies is None: logger.info( "domclick_city_sweep run_id=%d: куки Sber ID сессии отсутствуют/протухли — " "свип идёт без инъекции (см. #3264)", run_id, ) await run_domclick_city_sweep( db, run_id=run_id, config=ctx.config, matcher=ctx.matcher, shutdown_requested=ctx.shutdown_requested, proxy_provider=ctx.proxy_provider, city_id=int(params.get("city_id", 4)), rooms=params.get("rooms"), pages=int(params.get("pages_per_anchor", 5)), request_delay_sec=float(params.get("request_delay_sec", 6.0)), resume_run_id=kit_pick_resume(db, run_id), cookies=cookies, ) # ── house_coords_from_listings — sync set-based UPDATE в executor (#2771) ───── async def _job_house_coords_from_listings( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.backfill_house_coords_from_listings import run_house_coords_from_listings loop = asyncio.get_event_loop() await loop.run_in_executor( None, lambda: run_house_coords_from_listings(db, run_id=run_id, params=params), ) # ── geoportal_coords_backfill — sync local exact match в executor (#1967) ───── async def _job_geoportal_coords_backfill( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.backfill_listings_coords_geoportal import run_geoportal_coords_backfill loop = asyncio.get_event_loop() await loop.run_in_executor( None, lambda: run_geoportal_coords_backfill(db, run_id=run_id, params=params), ) # ── cadastral_geo_match — sync FDW-refresh+match в executor ─────────────────── async def _job_cadastral_geo_match( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.cadastral_geo_match import run_cadastral_geo_match loop = asyncio.get_event_loop() await loop.run_in_executor( None, lambda: run_cadastral_geo_match(db, run_id=run_id, params=params), ) # ── osm_poi_ekb_refresh — sync FDW-refresh в executor ───────────────────────── async def _job_osm_poi_ekb_refresh( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.osm_poi_ekb_refresh import run_osm_poi_ekb_refresh loop = asyncio.get_event_loop() await loop.run_in_executor( None, lambda: run_osm_poi_ekb_refresh(db, run_id=run_id, params=params), ) # ── dtp_stat_refresh — sync ZIP-скачивание+парс ДТП в executor ──────────────── async def _job_dtp_stat_refresh( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.dtp_stat_refresh import run_dtp_stat_refresh loop = asyncio.get_event_loop() await loop.run_in_executor( None, lambda: run_dtp_stat_refresh(db, run_id=run_id, params=params), ) # ── house_imv_backfill — async Avito-IMV с heartbeat, lifecycle в job ───────── async def _job_house_imv_backfill( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: batch_size = int(params.get("batch_size", 50)) request_delay_sec = float(params.get("request_delay_sec", 5.0)) only_status = str(params.get("only_status", "pending")) def _heartbeat() -> None: # heartbeat каждые N домов — иначе reap_zombies снимет живой долгий IMV-run (#1363). ctx.runs.update_heartbeat(db, run_id, {}) try: result = await ctx.enrichment.house_imv_backfill( db, batch_size=batch_size, request_delay_sec=request_delay_sec, only_status=only_status, heartbeat=_heartbeat, ) counters = { "checked": result.checked, "saved": result.saved, "skipped": result.skipped, "errors": result.errors, "duration_sec": int(result.duration_sec), # #2674: _column_counts (scrape_runs.py) берёт выделенные колонки из # ключей total_seen|lots_fetched и new_count|lots_inserted — ни одного # из них тут не было, поэтому все 39 прогонов этого source лежат в БД # с total_seen=0. А mark_done по этой же колонке шлёт алерт «3 подряд # done с нулевым результатом» (#2625) — то есть даже идеальный прогон # с 50 сохранёнными считался бы нулевым и через три дня выстрелил бы # ложной тревогой про капчу. # Трейд-офф: на исчерпанной очереди checked=0 три дня подряд тоже даст # алерт — но пустая очередь при ежедневном расписании это и правда сигнал. "total_seen": result.checked, "new_count": result.saved, # #2674: из скольких слотов пакета взяты дома на ПОВТОР (transient_error) # и сколько домов ушло в no_params до пакета одним запросом. Без этих # двух счётчиков в scrape_runs.counters проверить, что застрявшие # действительно возвращаются в очередь, можно только по houses. "retried": result.retried, "premarked": result.premarked, } # Честный статус (#2674, тот же класс, что #2670/#2657): успех — это # «сделали то, что собирались», а не «не поймали известное исключение». # На проде так ушли в done 31 прогон подряд: saved=0 при errors≈35 из 50. # Ноль сохранённых БЕЗ ошибок (всё отфильтровано в skipped) — честная # пустота, она по-прежнему done. if result.saved == 0 and result.errors > 0: ctx.runs.mark_failed( db, run_id, f"saved=0 при errors={result.errors} (checked={result.checked})", counters, ) else: ctx.runs.mark_done(db, run_id, counters) except Exception as exc: logger.exception("scheduler: house_imv_backfill crashed run_id=%d", run_id) try: ctx.runs.mark_failed(db, run_id, str(exc)[:1000], {}) except Exception: logger.exception("scheduler: mark_failed crashed run_id=%d", run_id) # ── domrf_kapremont_load — sync загрузка open data ДОМ.РФ в executor ───────── # #2674: loader (services/domrf_kapremont_loader.py) и CLI (tasks/domrf_kapremont_load.py) # написаны и покрыты тестами с #2013, но Handler'а и строки расписания не было — источник # запускали руками ровно один раз, 12.07.2026 (29 978 строк, один и тот же loaded_at у всех). # Это не мёртвый код, а оборванная проводка: нечему было его вызвать. async def _job_domrf_kapremont_load( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: """Скачать КР1.1+КР1.2 ДОМ.РФ → staging → backfill houses → propagate listings. Тело переиспользует те же три функции, что и CLI (дизайн-инвариант модуля: не дублируем логику). Lifecycle не свой — mark_done/mark_failed здесь, как у _job_yandex_newbuilding_sweep. Счётчики кладём в total_seen/new_count: `scrape_runs._column_counts` берёт выделенные колонки именно из этих ключей, и по ним же mark_done ловит «три подряд нулевых прогона» (#2625) — без них идеальный прогон лежал бы в БД как нулевой (тот же промах, что чинили у house_imv_backfill). """ from app.services.domrf_kapremont_loader import ( backfill_houses_from_domrf, load_domrf_kapremont, propagate_listings_year_from_houses, ) def _run() -> dict[str, int]: load_counts = load_domrf_kapremont(db) db.commit() houses_counts = backfill_houses_from_domrf(db) listings_counts = propagate_listings_year_from_houses(db) db.commit() return { "kr11_rows": load_counts["kr11_rows"], "upserted": load_counts["upserted"], "houses_updated": houses_counts["houses_updated"], "listings_updated": listings_counts["listings_updated"], # см. докстринг: выделенные колонки прогона + гейт «нулевой прогон». "total_seen": load_counts["kr11_rows"], "new_count": houses_counts["houses_updated"] + listings_counts["listings_updated"], } loop = asyncio.get_event_loop() try: counters = await loop.run_in_executor(None, _run) ctx.runs.mark_done(db, run_id, counters) except Exception as exc: logger.exception("scheduler: domrf_kapremont_load crashed run_id=%d", run_id) db.rollback() ctx.runs.mark_failed(db, run_id, str(exc)[:1000], {}) # ── fns_opendata_load — bulk-дампы ФНС по юрлицам → fns_legal_entity_facts ──── # Открытые данные ФНС по ЮРЛИЦАМ (лицензия nalog.gov.ru/opendata разрешает # переработку/перераспространение) — легальный обход того, что открытых данных # ЕГРН по правообладателям-физлицам не существует (218-ФЗ ст. 62). Потребителя у # fns_legal_entity_facts на момент добавления НЕТ (см. docstring # app/services/fns_opendata_loader.py и fns_lookup.py) — расписание сидируется # enabled=false, включение отдельным осознанным шагом. async def _job_fns_opendata_load( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: """Скачать наборы ФНС opendata (revexp/sshr2019/debtam/snr) → fns_legal_entity_facts. Тело переиспользует load_dataset (тот же дизайн-инвариант модуля, что у _job_domrf_kapremont_load: не дублируем логику CLI). commit() после каждого набора — сбой на debtam не должен откатывать уже загруженный revexp. """ from app.services.fns_opendata_loader import DATASET_SLUGS, load_dataset datasets = params.get("datasets") or list(DATASET_SLUGS) def _run() -> dict[str, int]: total_records = 0 total_upserted = 0 for slug in datasets: result = load_dataset(db, slug) db.commit() total_records += int(result["records"]) total_upserted += int(result["upserted"]) return { "records": total_records, "upserted": total_upserted, # см. докстринг _job_domrf_kapremont_load: выделенные колонки прогона + # гейт «три подряд нулевых прогона». "total_seen": total_records, "new_count": total_upserted, } loop = asyncio.get_event_loop() try: counters = await loop.run_in_executor(None, _run) ctx.runs.mark_done(db, run_id, counters) except Exception as exc: logger.exception("scheduler: fns_opendata_load crashed run_id=%d", run_id) db.rollback() ctx.runs.mark_failed(db, run_id, str(exc)[:1000], {}) # ── frt_mkd_load — АИС ППК ФРТ, реестр МКД region 66 (issue #frt-mkd) ──────── async def _job_frt_mkd_load( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: """Скачать реестр МКД АИС ФРТ (node_id) → staging frt_mkd → backfill houses. Тело переиспользует те же функции, что и CLI (app/tasks/frt_mkd_load.py) — не дублируем логику. node_id/region_code берутся из scrape_schedules.default_params (мигр. 291), с фоллбеком на дефолты loader'а. """ from app.services.frt_mkd_loader import ( DEFAULT_NODE_ID, DEFAULT_REGION_CODE, backfill_houses, load_frt_mkd, ) node_id = int(params.get("node_id", DEFAULT_NODE_ID)) region_code = int(params.get("region_code", DEFAULT_REGION_CODE)) def _run() -> dict[str, int]: load_counts = load_frt_mkd(db, node_id=node_id, region_code=region_code) db.commit() houses_counts = backfill_houses(db) db.commit() return { "rows": load_counts["rows"], "upserted": load_counts["upserted"], "houses_updated": houses_counts["houses_updated"], "total_seen": load_counts["rows"], "new_count": houses_counts["houses_updated"], } loop = asyncio.get_event_loop() try: counters = await loop.run_in_executor(None, _run) ctx.runs.mark_done(db, run_id, counters) except Exception as exc: logger.exception("scheduler: frt_mkd_load crashed run_id=%d", run_id) db.rollback() ctx.runs.mark_failed(db, run_id, str(exc)[:1000], {}) # ── cbr_macro_pull — макро-ряды ЦБ РФ (ипотека по субъектам + ключевая ставка) ─ async def _job_cbr_macro_pull( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: """Ипотечные ряды (XLSX) + ключевая ставка (SOAP) → cbr_mortgage_series/cbr_key_rate. Тело переиспользует те же функции, что и CLI (app/tasks/cbr_macro_pull.py) — дизайн-инвариант модуля: не дублируем логику. Lifecycle не свой — mark_done/mark_failed здесь, как у _job_domrf_kapremont_load. """ from app.services.cbr_macro import load_key_rate, pull_cbr_mortgage def _run() -> dict[str, int]: mortgage_counts = pull_cbr_mortgage(db) db.commit() key_rate_counts = load_key_rate(db) db.commit() return { "mortgage_upserted": mortgage_counts["upserted"], "mortgage_skipped": mortgage_counts["skipped"], "mortgage_errors": mortgage_counts["errors"], "key_rate_upserted": key_rate_counts["upserted"], # выделенные колонки прогона + гейт «три подряд нулевых прогона» (#2625). "total_seen": mortgage_counts["upserted"] + key_rate_counts["rows"], "new_count": mortgage_counts["upserted"] + key_rate_counts["upserted"], } loop = asyncio.get_event_loop() try: counters = await loop.run_in_executor(None, _run) ctx.runs.mark_done(db, run_id, counters) except Exception as exc: logger.exception("scheduler: cbr_macro_pull crashed run_id=%d", run_id) db.rollback() ctx.runs.mark_failed(db, run_id, str(exc)[:1000], {}) # ── purge_expired_trade_in_data — ЭТАП 4 B2C retention (152-ФЗ) ─────────────── async def _job_purge_expired_trade_in_data( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.tasks.purge_expired_trade_in_data import purge_expired_trade_in_data batch_size = params.get("batch_size") max_batches = params.get("max_batches") loop = asyncio.get_event_loop() await loop.run_in_executor( None, lambda: purge_expired_trade_in_data( db, run_id, batch_size=batch_size, max_batches=max_batches ), ) # ── house_dedup_merge — sync destructive merge в executor, owns lifecycle ───── async def _job_house_dedup_merge( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.services.house_dedup_merge import run_house_dedup_merge loop = asyncio.get_event_loop() await loop.run_in_executor( None, lambda: run_house_dedup_merge(db, run_id=run_id, params=params), ) # ── proxy_healthcheck — async, sub-hourly reschedule + lifecycle в job ──────── async def _job_proxy_healthcheck( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.services.proxy_pool import run_proxy_healthcheck try: counters = await run_proxy_healthcheck(db) ctx.runs.mark_done(db, run_id, counters) except Exception as exc: logger.exception("scheduler: run_proxy_healthcheck crashed run_id=%d", run_id) try: ctx.runs.mark_failed(db, run_id, str(exc)[:1000], {}) except Exception: logger.exception("scheduler: mark_failed crashed run_id=%d", run_id) def build_product_handlers(ctx: SchedulerContext) -> dict[str, Handler]: """Реестр НЕ-sweep продуктовых source→Handler для kit build_registry. Kit-native sweeps (avito/yandex/cian/domclick city/full-load/newbuilding) НЕ здесь — их даёт build_registry(_default_kit_handlers). Здесь — именованные + 1 wildcard (deactivate_stale_*), покрывающие каждый НЕ-sweep source боевого scheduler-dispatch. (Число намеренно не названо: прежнее «19» разошлось с реальностью на пять записей.) `ctx` — принят для симметрии контракта; сами Handler-job'ы получают ctx во время dispatch (см. kit `_dispatch`), поэтому здесь он не замыкается. """ return { "cian_history_backfill": Handler( _job_cian_history_backfill, "cian_history_backfill", pre_claim=_cian_pre_claim, ), # #3284: то же тело, другая очередь. Разные source нужны именно как РАЗНЫЕ # строки расписания — у них свои окна, свой next_run_at и свой счётчик # прогонов; гонять оба режима под одним source нельзя, планировщик держит # на source ровно один активный прогон. # Каденс — в минутах (замер прода 31.08.2026): та же суточная дыра # compute_next_run_at, что у avito/yandex_detail_backfill — Циан держит # темп ~23 с/карточку при очереди 20 501 и НОЛЕ блоков за неделю. # # 360 минут (4 прогона/сутки), НЕ 180 как у Яндекса — АСИММЕТРИЯ ниже # ОБЯЗАТЕЛЬНА к прочтению перед тем, как выравнивать эти два интервала: # Циан ходит через BrowserFetcher и ДЕРЖИТ lease узла на ВЕСЬ прогон, а # не только берёт URL прокси (как Яндекс через resolve_proxy_url/ # proxy_egress). Пул сейчас 3 живых узла — дефицитный ресурс, который # Циан на время прогона отбирает у avito_detail_backfill и # domclick_detail_backfill. Учащение Циана до 180 мин без замера того, # как это скажется на соседях по пулу, — прямой риск сжечь мощность, # которой и так не хватает на троих. "cian_detail_backfill": Handler( _job_cian_history_backfill, "cian_detail_backfill", pre_claim=_cian_pre_claim, post_claim=reschedule_after_minutes(param="interval_minutes", default=360), ), "rosreestr_dkp_import": Handler(_job_rosreestr_dkp, "rosreestr_dkp_import"), # Wildcard (#3051 п.3): rosreestr_dkp_import_77 (Москва, миграция 288) и любой # будущий region_code-суффикс из той же семьи резолвятся сюда через # resolve_handler по префиксу (тот же механизм, что deactivate_stale_* / # avito_city_sweep_* — см. scraper_kit.orchestration.scheduler.resolve_handler). # region_code берётся из default_params строки расписания (import_rosreestr_dkp # сам валидирует его через app.services.regions.REGIONS), Handler-тело общее. "rosreestr_dkp_import_*": Handler(_job_rosreestr_dkp, "rosreestr_dkp_import_*"), "listing_source_snapshot": Handler(_job_listing_source_snapshot, "listing_source_snapshot"), "asking_to_sold_ratio_refresh": Handler( _job_asking_to_sold_ratio, "asking_to_sold_ratio_refresh" ), "deal_city_price_bands_refresh": Handler( _job_deal_city_price_bands_refresh, "deal_city_price_bands_refresh" ), "refresh_search_matview": Handler(_job_refresh_search_matview, "refresh_search_matview"), "yandex_address_backfill": Handler(_job_yandex_address_backfill, "yandex_address_backfill"), "sber_index_pull": Handler(_job_sber_index_pull, "sber_index_pull"), "rosreestr_quarter_poll": Handler(_job_rosreestr_quarter_poll, "rosreestr_quarter_poll"), "deals_freshness_monitor": Handler(_job_deals_freshness_monitor, "deals_freshness_monitor"), "sber_freshness_monitor": Handler(_job_sber_freshness_monitor, "sber_freshness_monitor"), "landing_stats_refresh": Handler(_job_landing_stats, "landing_stats_refresh"), "landing_showcase_deals": Handler(_job_landing_showcase_deals, "landing_showcase_deals"), "newbuilding_enrich": Handler(_job_newbuilding_enrich, "newbuilding_enrich"), "yandex_newbuilding_sweep": Handler( _job_yandex_newbuilding_sweep, "yandex_newbuilding_sweep" ), "geoportal_coords_backfill": Handler( _job_geoportal_coords_backfill, "geoportal_coords_backfill" ), "house_coords_from_listings": Handler( _job_house_coords_from_listings, "house_coords_from_listings" ), "geocode_missing_listings": Handler( _job_geocode_missing_listings, "geocode_missing_listings" ), # Каденс — в минутах, а не в сутках (замер 2026-08-22). Дефолтная # гранулярность compute_next_run_at — сутки, и бэкфилл получал ровно один # прогон в день. При этом прогон умирает по бану через 17-83 минуты, то # есть 23 часа из 24 задание простаивало: 238 обогащённых карточек за # сутки при 9 951 активном объявлении — очередь разбиралась бы месяцами. # # Хук сам себя throttl'ит: пока прогон идёт, has_running_run в _claim_run # возвращает None и next_run_at не сбрасывается. # # 180 минут — ОСОЗНАННО консервативная отправная точка, а не найденный # оптимум. Данных для подбора нет: два прогона дали противоречивую # картину (4562 — 83 мин и 175 карточек; 4586 через 4.3 часа, когда пул # прокси был давно чист, — 17 мин и 42 карточки). Значит память Авито # длиннее часов, и учащение может ухудшить выход, а не улучшить. # # Риск, который надо держать в голове при подборе: те же 4 прокси # обслуживают SERP-свипы — первичный сбор. Сжечь их на обогащении хуже, # чем медленно обогащать. Двигать интервал вниз только по замеру # нескольких суток, глядя и на свипы тоже. "avito_detail_backfill": Handler( _job_avito_detail_backfill, "avito_detail_backfill", post_claim=reschedule_after_minutes(param="interval_minutes", default=180), ), # Каденс — в минутах (замер прода 31.08.2026), тот же дефект, что у # avito_detail_backfill выше: без post_claim compute_next_run_at держит # суточную гранулярность, а прогон Яндекса разбирает 375-450 карточек в # час (budget_sec=3600) при очереди 11 110 и НУЛЕ блоков за неделю — узкое # место чисто в частоте запуска, не в самом фетчере. # # 180 минут (8 прогонов/сутки) безопасны именно из-за асимметрии с Циан: # Яндекс ходит через resolve_proxy_url (proxy_egress) — берёт URL прокси, # но НЕ держит lease узла на весь прогон, поэтому учащение не отнимает # узел у других source. "yandex_detail_backfill": Handler( _job_yandex_detail_backfill, "yandex_detail_backfill", post_claim=reschedule_after_minutes(param="interval_minutes", default=180), ), "domclick_detail_backfill": Handler( _job_domclick_detail_backfill, "domclick_detail_backfill" ), # Override kit-native (#3264) — инъекция куки сессии, см. докстринг job'а выше. "domclick_city_sweep": Handler(_job_domclick_city_sweep, "domclick_city_sweep"), "cadastral_geo_match": Handler(_job_cadastral_geo_match, "cadastral_geo_match"), "osm_poi_ekb_refresh": Handler(_job_osm_poi_ekb_refresh, "osm_poi_ekb_refresh"), "dtp_stat_refresh": Handler(_job_dtp_stat_refresh, "dtp_stat_refresh"), "house_imv_backfill": Handler(_job_house_imv_backfill, "house_imv_backfill"), "house_dedup_merge": Handler(_job_house_dedup_merge, "house_dedup_merge"), "domrf_kapremont_load": Handler(_job_domrf_kapremont_load, "domrf_kapremont_load"), "frt_mkd_load": Handler(_job_frt_mkd_load, "frt_mkd_load"), "cbr_macro_pull": Handler(_job_cbr_macro_pull, "cbr_macro_pull"), "fns_opendata_load": Handler(_job_fns_opendata_load, "fns_opendata_load"), "purge_expired_trade_in_data": Handler( _job_purge_expired_trade_in_data, "purge_expired_trade_in_data" ), "proxy_healthcheck": Handler( _job_proxy_healthcheck, "proxy_healthcheck", post_claim=reschedule_after_minutes(param="interval_minutes", default=30), ), # Wildcard: deactivate_stale_avito / _yandex / _cian — один Handler на семейство # (resolve_handler матчит по префиксу key[:-1]="deactivate_stale_"). "deactivate_stale_*": Handler(_job_deactivate_stale, "deactivate_stale"), } __all__ = ["build_product_handlers"]