All checks were successful
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Successful in 9s
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m50s
Deploy Trade-In / build-backend (push) Successful in 1m41s
Deploy Trade-In / deploy (push) Successful in 1m42s
683 lines
35 KiB
Python
683 lines
35 KiB
Python
"""Продуктовые 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"]
|