gendesign/tradein-mvp/backend/app/services/product_handlers.py
bot-backend b800760c24
All checks were successful
CI Trade-In / changes (pull_request) Successful in 7s
CI / changes (pull_request) Successful in 7s
CI / backend-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Successful in 1m3s
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 2m46s
fix(tradein/scraper): фильтр skipped в админке + освежение схлопнутой строки (#2658)
Правки по ревью PR #2662.

Фильтр статуса. `GET /admin/scrape/runs?status=skipped` отдавал 422 — 'skipped' не
было в Literal, а во фронте не было чипа. Строки рисовались, но задать вопрос
«что сейчас пропускается» на единственной поверхности, построенной ровно для
этого, было нельзя. Добавлено в оба места (translateStatus «пропущено» и
нейтральный бейдж уже умели).

Схлопывание освежает строку. UPDATE двигал только finished_at/heartbeat_at, из-за
чего живой стрик замерзал: списки прогонов сортируют ORDER BY started_at DESC и
берут limit=20, поэтому 37-дневный пропуск утонул бы под свежими прогонами других
источников — след в базе есть, на экране нет. Теперь started_at = NOW(), а начало
стрика переезжает в counters.first_skip_at; сортировку общего списка не трогаем
(она про все источники, чинить надо было одну строку). Там же обновляется
counters.detail — иначе в строке 37 дней висел текст «протухли 1 день назад»,
хотя именно эта цифра и есть предмет issue. jsonb_set заменён на `||` +
jsonb_build_object: три вложенных jsonb_set читать в 3 ночи невозможно, а NULL в
jsonb_set обнуляет весь counters.

Поиск последней строки. `ORDER BY id DESC` не ложится на индекс
(source, started_at DESC) из миграции 015 — для unknown_source (тикает каждые
60 с бессрочно) это отбор всех строк источника с сортировкой раз в минуту.
Теперь ORDER BY started_at DESC, id DESC.

session_expires_at получил valid_only: предупреждение «скоро протухнут» считает
срок ИМЕННО той записи, которую взял load_session — при нескольких аккаунтах
свежайшая-любая может быть чужой протухшей строкой. Диагностика после None
по-прежнему смотрит на свежайшую любую (валидных там нет по определению).

Запись пропуска намеренно НЕ обёрнута в свой try/except: если db.execute падает,
то падает и claim следующего расписания в этом же тике — тик срывается в любом
случае, а глушить исключение здесь значило бы вернуть ровно тот немой пропуск,
ради которого заведён #2658. Самовосстановление через 60 с.
2026-08-05 23:10:47 +05:00

513 lines
24 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Продуктовые 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 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")
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,
),
)
# ── 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")),
)
ctx.runs.mark_done(db, run_id, result.to_dict())
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)
# ── 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,
)
ctx.runs.mark_done(
db,
run_id,
{
"checked": result.checked,
"saved": result.saved,
"skipped": result.skipped,
"errors": result.errors,
"duration_sec": int(result.duration_sec),
},
)
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)
# ── 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). Здесь — 19 именованных + 1 wildcard
(deactivate_stale_*), покрывающие каждый НЕ-sweep source боевого scheduler-dispatch.
`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"
),
"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"),
"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"]