Топология подтверждена перед удалением (docker-compose.prod.yml): tradein-backend (uvicorn app.main:app) — SCHEDULER_ENABLE=false; tradein-scraper (python -m app.scheduler_main) — SCHEDULER_ENABLE=true + USE_KIT_SCHEDULER=true. Kit-путь (_run_kit_scheduler → scraper_kit.orchestration.scheduler + product_handlers) самодостаточен: не импортирует ничего из app.services.scheduler.scheduler_loop или app.services.scrape_pipeline. Все НЕ-sweep джобы, которые kit-scheduler диспетчерит через build_product_handlers, идут напрямую в app.tasks.*/ app.services.* (либо lazy-импортят import_rosreestr_dkp/_execute_cian_backfill из scheduler.py) — мимо удаляемой legacy-машинерии. app/services/scheduler.py: 2098 → 418 строк. Удалено: scheduler_loop, get_due_schedules, reap_zombies, _claim_run, _defer_next_run_at, _spawn_tracked/ _drain_inflight/_inflight_tasks, все 27 trigger_*_run-функций, импорт app.services.scrape_pipeline, константы SCHEDULER_TICK_SEC/ZOMBIE_THRESHOLD_HOURS (достижимы были только через удалённый scheduler_loop-путь). Оставлено (живые импортёры вне удалённого): compute_next_run_at + has_running_run (admin.py), import_rosreestr_dkp + _execute_cian_backfill (lazy-импорты в product_handlers.py — job-тела kit-handler'ов). main.py: убран `from app.services.scheduler import scheduler_loop` + lifespan-блок запуска (`if settings.scheduler_enable: asyncio.create_task(scheduler_loop())`); прод-backend всегда шёл с SCHEDULER_ENABLE=false, так что это был мёртвый код. scheduler_main.py: убрана ship-dark развилка #2192 (USE_KIT_SCHEDULER=false → legacy scheduler_loop fallback) — _run_kit_scheduler() теперь безусловный путь. Поле settings.use_kit_scheduler оставлено в конфиге (Settings extra="ignore" защищает от startup-краха на leftover env var), но на ветвление не влияет. app.services.scrape_pipeline: 0 runtime-импортёров в app/+scripts/+packages/ после этого PR (только тесты, которые Part E удалит вместе с самим файлом) — подтверждено grep. scrape_pipeline.py не тронут (Part E). Тесты: удалены test_house_imv_backfill_scheduler.py (100% legacy-триггер, backfill_house_imv сервис покрыт в test_house_imv_backfill_browser_flag.py / test_backfill_wave2.py) и test_kit_registry_completeness.py (parity-инвариант против удалённого dispatch, дублирует test_scraper_kit_scheduler_parity.py). Точечно вырезаны "Scheduler wiring" секции (trigger_fn_exists/dispatch_branch_ wired/runs_in_executor) из ~10 файлов, тестирующих сами task-функции — сами task-тесты (SQL-shape, миграции, fake-db поведение) оставлены нетронутыми. test_scheduler.py: 825 → ~90 строк (остались только compute_next_run_at-тесты). test_scraper_kit_scheduler_parity.py: убрана golden-parity секция против удалённого scheduler_loop (SOURCE_TO_OLD_TRIGGER/_drive_old_one_tick/ test_routing_parity_per_source), остальное (claim/reap_zombies/dispatch/ registry-shape тесты kit-модуля) сохранено — источник этих инвариантов не app.services.scheduler, а сам scraper_kit.orchestration.scheduler. test_scheduler_main.py: 2 теста, патчившие app.services.scheduler.scheduler_loop, переведены на монкипатч sm._run_kit_scheduler (единственный путь после этого PR). test_sweep_imv_phase.py:171-371 (6 прямых импортов run_avito_city_sweep из scrape_pipeline) намеренно НЕ тронуты — Part E. Verify: полный pytest 3179 passed / 6 skipped / 1 known-unrelated fail (test_search_cache_hit, #2208, не связан с этим PR); ruff 0.7.4 чист на всех изменённых файлах; `python -c "import app.main; import app.scheduler_main"` OK.
186 lines
9.7 KiB
Python
186 lines
9.7 KiB
Python
"""Deactivate stale listings that have not been seen in >TTL days.
|
||
|
||
Первоначально создано для avito (#759). Теперь поддерживает любой источник
|
||
(listing_source) с опциональной фильтрацией по listing_segment.
|
||
|
||
Ключевые решения:
|
||
- Cian/Yandex не поддерживают full-coverage sweep -> паушальный TTL сломает живой
|
||
инвентарь. DECISION: для yandex/cian деактивировать ТОЛЬКО listing_segment='vtorichka',
|
||
TTL=30. novostroyki (9659 активных первичных строк) и NULL-сегмент не трогаем.
|
||
- avito: все сегменты (segments=None), TTL=10 дней -- поведение без изменений.
|
||
- Строки НЕ удаляются -- история нужна для бэктеста (#667).
|
||
|
||
Задача синхронная (DB-only, никаких внешних HTTP-вызовов) -- запускается kit-scheduler'ом
|
||
через product_handlers._job_deactivate_stale (wildcard-handler deactivate_stale_*),
|
||
по образцу snapshot_listing_sources / recompute_asking_to_sold_ratios
|
||
(sync task в run_in_executor).
|
||
|
||
TTL для avito берётся из settings.avito_stale_ttl_days (env AVITO_STALE_TTL_DAYS, default 10).
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
from typing import Any
|
||
|
||
from sqlalchemy import text
|
||
from sqlalchemy.orm import Session
|
||
|
||
from app.core.config import settings
|
||
from app.services import scrape_runs as runs_mod
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# Разрешённые колонки-таймстемпы для проверки свежести (staleness).
|
||
# Whitelist ЖЁСТКО ограничивает f-string-интерполяцию имени колонки в SQL:
|
||
# значение колонки НИКОГДА не берётся из пользовательского ввода без этой проверки.
|
||
# - last_seen_at: дефолт; двигается при каждом «увиденном» объявлении в прогоне.
|
||
# - scraped_at: честная свежесть для источников с нетрекаемым bulk-touch
|
||
# last_seen_at (domklik: #2204 — bulk UPDATE двигает last_seen_at всем строкам
|
||
# одним timestamp, поэтому TTL по last_seen_at бесполезен; scraped_at двигает
|
||
# только реальный скрейп).
|
||
_ALLOWED_STALENESS_COLUMNS = frozenset({"last_seen_at", "scraped_at"})
|
||
|
||
|
||
def _build_all_segments_sql(staleness_column: str) -> Any:
|
||
"""UPDATE без фильтра по сегменту: все сегменты для данного source.
|
||
|
||
staleness_column уже прошёл whitelist-проверку в deactivate_stale_listings,
|
||
поэтому f-string-подстановка имени колонки безопасна. Значения (:listing_source,
|
||
:ttl_days) остаются param-binding — psycopg v3 safe (никаких :param::type).
|
||
"""
|
||
return text(
|
||
f"""
|
||
UPDATE listings
|
||
SET is_active = false
|
||
WHERE source = :listing_source
|
||
AND is_active = true
|
||
AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval)
|
||
"""
|
||
)
|
||
|
||
|
||
def _build_segments_sql(staleness_column: str) -> Any:
|
||
"""UPDATE с фильтром по сегменту (segments задан): только указанные сегменты.
|
||
|
||
staleness_column уже прошёл whitelist-проверку. = ANY(CAST(:segments AS text[]))
|
||
-- psycopg v3 адаптирует Python list -> text[].
|
||
"""
|
||
return text(
|
||
f"""
|
||
UPDATE listings
|
||
SET is_active = false
|
||
WHERE source = :listing_source
|
||
AND is_active = true
|
||
AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval)
|
||
AND listing_segment = ANY(CAST(:segments AS text[]))
|
||
"""
|
||
)
|
||
|
||
|
||
# Дефолтные (last_seen_at) варианты SQL -- сохранены как модульные константы для
|
||
# обратной совместимости (тесты читают .text, product_handlers/scheduler не менялись).
|
||
_DEACTIVATE_SQL_ALL_SEGMENTS = _build_all_segments_sql("last_seen_at")
|
||
_DEACTIVATE_SQL_SEGMENTS = _build_segments_sql("last_seen_at")
|
||
|
||
# Алиас для обратной совместимости -- использовался в тестах через task_mod._DEACTIVATE_SQL
|
||
_DEACTIVATE_SQL = _DEACTIVATE_SQL_ALL_SEGMENTS
|
||
|
||
|
||
def deactivate_stale_listings(
|
||
db: Session,
|
||
run_id: int,
|
||
*,
|
||
listing_source: str,
|
||
ttl_days: int,
|
||
segments: list[str] | None = None,
|
||
staleness_column: str = "last_seen_at",
|
||
) -> dict[str, int]:
|
||
"""Пометить is_active=false объявления, чья свежесть старше ttl_days дней.
|
||
|
||
Параметры:
|
||
listing_source: значение поля source в таблице listings ('avito', 'yandex',
|
||
'cian', 'domklik').
|
||
ttl_days: количество дней TTL; объявления старше этого порога деактивируются.
|
||
segments: если задан -- деактивировать только объявления с указанными
|
||
listing_segment значениями. None -> все сегменты (поведение avito по умолчанию).
|
||
staleness_column: колонка-таймстемп, по которой считается свежесть. Whitelist
|
||
{"last_seen_at", "scraped_at"} — иначе ValueError ДО любого SQL. Дефолт
|
||
last_seen_at. Для domklik (#2204) — scraped_at: нетрекаемый bulk-touch
|
||
двигает last_seen_at всем строкам одним timestamp, поэтому честная
|
||
свежесть = scraped_at (двигается только реальным скрейпом).
|
||
|
||
Sync (вызывается scheduler-триггером в executor, как snapshot_listing_sources).
|
||
Один UPDATE в транзакции. Финализирует scrape_runs (mark_done / mark_failed).
|
||
|
||
Returns {"deactivated": N} -- количество обновлённых строк.
|
||
|
||
Raises:
|
||
ValueError: если staleness_column не входит в whitelist (проверка ДО SQL,
|
||
никакой интерполяции пользовательского ввода в запрос).
|
||
"""
|
||
counters: dict[str, int] = {"deactivated": 0}
|
||
try:
|
||
# Whitelist-проверка ДО построения/выполнения SQL: только после неё имя колонки
|
||
# интерполируется f-string'ом. Значения по-прежнему идут через param-binding.
|
||
# Внутри try -> невалидная колонка финализирует run как failed (mark_failed),
|
||
# а не оставляет его «running» (единый контракт с прочими сбоями).
|
||
if staleness_column not in _ALLOWED_STALENESS_COLUMNS:
|
||
raise ValueError(
|
||
f"invalid staleness_column={staleness_column!r}; "
|
||
f"allowed: {sorted(_ALLOWED_STALENESS_COLUMNS)}"
|
||
)
|
||
|
||
# segments is None -> все сегменты (поведение avito). segments=[...] -> только
|
||
# перечисленные сегменты. Используем `is not None` (НЕ truthy): пустой список []
|
||
# означает "ни один сегмент" (= ANY(ARRAY[]) ничего не матчит, деактивирует 0),
|
||
# а НЕ "все сегменты" — иначе случайный [] стёр бы весь источник.
|
||
if segments is not None:
|
||
params: dict[str, Any] = {
|
||
"listing_source": listing_source,
|
||
"ttl_days": ttl_days,
|
||
"segments": segments,
|
||
}
|
||
result = db.execute(_build_segments_sql(staleness_column), params)
|
||
else:
|
||
params = {"listing_source": listing_source, "ttl_days": ttl_days}
|
||
result = db.execute(_build_all_segments_sql(staleness_column), params)
|
||
|
||
counters["deactivated"] = result.rowcount or 0
|
||
|
||
db.commit()
|
||
runs_mod.mark_done(db, run_id, counters)
|
||
logger.info(
|
||
"deactivate_stale source=%s run_id=%d done: deactivated=%d "
|
||
"(ttl_days=%d, segments=%r, staleness_column=%s)",
|
||
listing_source,
|
||
run_id,
|
||
counters["deactivated"],
|
||
ttl_days,
|
||
segments,
|
||
staleness_column,
|
||
)
|
||
return counters
|
||
except Exception as exc:
|
||
logger.exception("deactivate_stale source=%s run_id=%d failed", listing_source, run_id)
|
||
db.rollback()
|
||
runs_mod.mark_failed(db, run_id, str(exc)[:1000], counters)
|
||
raise
|
||
|
||
|
||
def deactivate_stale_avito_listings(db: Session, run_id: int) -> dict[str, int]:
|
||
"""Обёртка для обратной совместимости -- deactivate_stale_avito (все сегменты, TTL из settings).
|
||
|
||
Пометить is_active=false все avito-объявления с last_seen_at > TTL дней назад.
|
||
Sync (вызывается scheduler-триггером в executor, как snapshot_listing_sources).
|
||
Один UPDATE в транзакции. Финализирует scrape_runs (mark_done / mark_failed).
|
||
|
||
Returns {"deactivated": N} -- количество обновлённых строк.
|
||
"""
|
||
return deactivate_stale_listings(
|
||
db,
|
||
run_id,
|
||
listing_source="avito",
|
||
ttl_days=settings.avito_stale_ttl_days,
|
||
segments=None,
|
||
)
|