"""Продуктовые 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 city/full-load) регистрируются самим `build_registry` через `_default_kit_handlers` — здесь их НЕТ. Дизайн-инвариант: НЕ дублируем логику. Каждый `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, ) 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_MIN_CONFIRMATIONS, 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) 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, ), ) # ── 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) # ── 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 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) # ── 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), ) # ── 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], {}) # ── 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, ), "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"), "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" ), "avito_detail_backfill": Handler(_job_avito_detail_backfill, "avito_detail_backfill"), "yandex_detail_backfill": Handler(_job_yandex_detail_backfill, "yandex_detail_backfill"), "domclick_detail_backfill": Handler( _job_domclick_detail_backfill, "domclick_detail_backfill" ), "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"), "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"), "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"]