gendesign/tradein-mvp/backend/app/scheduler_main.py
bot-backend 1f85ef7d4e
All checks were successful
CI / changes (pull_request) Successful in 8s
CI Trade-In / backend-tests (pull_request) Successful in 3m50s
CI Trade-In / changes (pull_request) Successful in 8s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
fix(tradein/payments): тело нотификации не течёт в мониторинг и аудит, повторы банка не отбиваются лимитом
PR-D2 платёжного контура МЕРЫ — закрывает утечки до открытия публичных путей
(PR-D3/D4), сам ничего не открывает: _PUBLIC_PATHS (rbac.py), Caddyfile,
roles.yaml, auth_session.py не тронуты.

- sentry_scrub.py: новая scrub_payment_request_body — вырезает
  event.request.data целиком для /api/v1/trade-in/payments/* (sentry_sdk 2.64
  кладёт полное тело запроса в request.data, send_default_pii=False это НЕ
  гейтит — тот флаг управляет только куками). Плюс расширен _PII_KEYS:
  customer_email/customer_phone/pan/expdate/cardid/rebillid/token/terminalkey.
- main.py, scheduler_main.py, tgbot_main.py (все 3 точки инициализации
  sentry_sdk.init в проекте) — тот же обработчик проведён в ОБА канала,
  before_send и before_send_transaction. Мотивирующий инцидент: на соседнем
  продукте вчера закрыли только error-канал, transaction остался без
  обработчика вообще.
- ratelimit.py: точный путь notify — свой щедрый SlidingWindowLimiter
  (3000/60с per-IP, идиома support.py) вместо общего лимитера, но НЕ полное
  отключение — backstop против шторма запросов остаётся, подпись проверяется
  уже после разбора тела (PR-D3). Только notify, не checkout (тот с сессией).
- request_audit.py: notify — в audit skip-набор (defense-in-depth: middleware
  внешний относительно rbac_guard и читает сырой X-Authenticated-User —
  спуфнутый заголовок иначе писал бы фальшивые события с атрибуцией admin).
- smoke-mera-perimeter.sh: негативные проверки-канарейки — notify/checkout
  сейчас закрыты 404 (meraocenka.ru, Caddy не проксирует) и 401
  (gendsgn.ru, rbac ещё не открыл) с обеих сторон периметра.

Тесты: scrub на произвольной глубине + payment-path body-wipe, AST-разбор
(не substring — комментарии в этих же файлах сами упоминают
before_send_transaction) на проводку обоих каналов во всех точках
инициализации, 400 запросов notify без единого 429 + контроль что общий
лимитер по-прежнему активен на других путях, notify вне user_events даже со
спуфнутым X-Authenticated-User: admin.
2026-08-07 16:17:59 +03:00

232 lines
11 KiB
Python
Raw 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.

"""Standalone entrypoint для tradein-scraper контейнера (#1182).
Зачем: API-деплой (docker restart tradein-backend) не должен прерывать
бегущие scraper sweep'ы (avito/cian/rosreestr). Вынесен в отдельный
контейнер tradein-scraper с тем же backend-образом, но другой командой.
Запуск: python -m app.scheduler_main
Zombie-reap: reap_zombies() вызывается в начале каждого тика kit scheduler_loop
(см. scraper_kit/orchestration/scheduler.py), дублировать здесь не нужно.
"""
from __future__ import annotations
import asyncio
import logging
import os
import signal
import sys
from contextlib import suppress
from app.core.config import settings
from app.core.shutdown import request_shutdown, shutdown_requested, wait_for_shutdown
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)s %(name)s: %(message)s",
)
logger = logging.getLogger(__name__)
# Safety net: после запроса drain'а ждём добровольного выхода задачи не дольше
# этого порога, который ОБЯЗАН быть < docker stop_grace_period (120s). Если
# некооперативная задача не вышла за это время — hard-cancel, чтобы процесс не
# подвис до SIGKILL за пределами grace. Module-level → monkeypatch'абельно в тестах.
_DRAIN_TIMEOUT_S = 100.0
# Мониторинг ошибок — GlitchTip (Sentry-совместимый, #396).
# Только integrations без Starlette/FastAPI (у нас нет ASGI-приложения здесь).
# SqlalchemyIntegration + HttpxIntegration + LoggingIntegration покрывают scraper-слой.
if settings.glitchtip_dsn:
import sentry_sdk
from sentry_sdk.integrations.httpx import HttpxIntegration
from sentry_sdk.integrations.logging import LoggingIntegration
from sentry_sdk.integrations.sqlalchemy import SqlalchemyIntegration
from app.observability.sentry_scrub import scrub_payment_request_body, scrub_pii_event
def _before_send(event: object, hint: dict[str, object]) -> object:
"""PR-D2: этот процесс не держит ASGI-приложения (нет `request` в event
сегодня), но payments_confirm/payments_reconcile (PR-E, тот же
`tradein-scraper` контейнер) будут звать Т-Банк API отсюда — belt-and-
suspenders на случай, если платёжные данные когда-нибудь попадут в
`request`/`extra`. Тот же обработчик на оба канала ниже — см.
app/main.py._before_send (идентичный мотив, не дублировать без причины)."""
scrubbed = scrub_payment_request_body(event, hint) # type: ignore[arg-type]
if scrubbed is None:
return None
return scrub_pii_event(scrubbed, hint) # type: ignore[arg-type]
sentry_sdk.init(
dsn=settings.glitchtip_dsn,
environment=settings.environment,
release=os.getenv("GIT_SHA") or os.getenv("SENTRY_RELEASE") or "unknown",
traces_sample_rate=0.0,
send_default_pii=False,
before_send=_before_send,
before_send_transaction=_before_send,
integrations=[
SqlalchemyIntegration(),
HttpxIntegration(),
LoggingIntegration(level=logging.INFO, event_level=logging.ERROR),
],
)
logger.info("GlitchTip monitoring enabled (scheduler_main)")
def _should_run() -> bool:
"""Проверить kill-switch SCHEDULER_ENABLE перед запуском asyncio loop."""
return settings.scheduler_enable
async def _await_scheduler(task: asyncio.Task[None]) -> None:
"""Дождаться завершения scheduler-задачи с кооперативным SIGTERM-drain'ом.
- Нет shutdown: scheduler_loop бесконечен → задача никогда не завершается, ждём её
как есть (процесс просто работает).
- shutdown запрошен, задача ещё бежит: bounded `wait_for(_DRAIN_TIMEOUT_S)` —
кооперативная задача докоммитит текущий unit на ближайшем checkpoint'е и выйдет
сама. Если превысила grace → hard-cancel + suppress CancelledError, чтобы
некооперативная задача не подвесила процесс за пределами docker stop_grace_period.
"""
# Гонка «задача завершилась сама» против «пришёл SIGTERM»: отсчёт safety-net'а
# должен стартовать от МОМЕНТА запроса drain'а, а не от старта процесса.
shutdown_waiter = asyncio.create_task(wait_for_shutdown())
try:
await asyncio.wait({task, shutdown_waiter}, return_when=asyncio.FIRST_COMPLETED)
finally:
shutdown_waiter.cancel()
with suppress(asyncio.CancelledError):
await shutdown_waiter
if task.done():
# Задача вышла сама (drained на checkpoint'е; без сигнала сюда не попадаем).
# .result() пробрасывает реальное исключение, если scheduler_loop неожиданно
# упал (сохраняем прежнюю loud-crash семантику `await task`), а не глушит его.
task.result()
logger.info("scheduler_main: scheduler task exited cleanly")
return
logger.info(
"scheduler_main: SIGTERM-drain — waiting up to %.0fs for in-flight unit to commit",
_DRAIN_TIMEOUT_S,
)
try:
await asyncio.wait_for(task, timeout=_DRAIN_TIMEOUT_S)
logger.info("scheduler_main: scheduler drained and exited cleanly")
except TimeoutError:
logger.warning(
"scheduler_main: drain exceeded %.0fs grace — hard-cancelling scheduler task",
_DRAIN_TIMEOUT_S,
)
task.cancel()
with suppress(asyncio.CancelledError):
await task
async def _run_kit_scheduler() -> None:
"""Kit-путь (#2192): собрать SchedulerContext + registry и крутить kit scheduler_loop.
Инжектируем продуктовые адаптеры (RealScraperConfig/Matcher/Enrichment/SessionFactory) +
боевой runs-модуль (app.services.scrape_runs) + текущий shutdown_requested callable, чтобы
kit-loop делил ту же concurrency/lifecycle-семантику с боевым scheduler'ом. Продуктовые
НЕ-sweep source'ы приходят из build_product_handlers; kit-native sweeps — из build_registry.
SIGTERM-drain: kit scheduler_loop сам следит за ctx.shutdown_requested() и дренит
in-flight задачи (ctx.drain_inflight) — та же кооперативная семантика, что у боевого.
"""
from scraper_kit.orchestration.scheduler import (
SchedulerContext,
build_registry,
)
from scraper_kit.orchestration.scheduler import (
scheduler_loop as kit_scheduler_loop,
)
from app.services import scrape_runs as runs_mod
from app.services.product_handlers import build_product_handlers
from app.services.scraper_adapters import (
RealEnrichmentJobs,
RealMatcherAdapter,
RealProxyProvider,
RealScraperConfig,
RealSessionFactory,
)
# Proxy-пул (#2163 curl / #2164 P4 browser) инжектируется ТОЛЬКО когда включён любой
# из флагов. Оба дефолтно False (ship-dark) → proxy_provider=None → curl/browser берут
# прокси из env, прод не меняется. Один RealProxyProvider обслуживает оба пути (curl
# через curl_proxy_url, browser через BrowserFetcher → POST /fetch{proxy}).
use_pool = settings.use_proxy_pool_curl or settings.use_proxy_pool_browser
proxy_provider = RealProxyProvider() if use_pool else None
ctx = SchedulerContext(
config=RealScraperConfig(),
matcher=RealMatcherAdapter(),
enrichment=RealEnrichmentJobs(),
session_factory=RealSessionFactory(),
shutdown_requested=shutdown_requested,
runs=runs_mod,
proxy_provider=proxy_provider,
)
registry = build_registry(product_handlers=build_product_handlers(ctx))
logger.info("scheduler_main: kit-scheduler registry built (%d source handlers)", len(registry))
await kit_scheduler_loop(ctx, registry)
async def _run() -> None:
"""Async entrypoint: kit-scheduler с кооперативным SIGTERM/SIGINT-drain'ом.
SIGTERM/SIGINT → request_shutdown() (НЕ task.cancel()): бегущий scrape-unit
докоммитит текущую карточку и выйдет сам на ближайшем between-unit checkpoint'е.
Safety-net в _await_scheduler гарантирует выход в пределах docker grace.
#2397 Part C: legacy app.services.scheduler.scheduler_loop fallback (USE_KIT_SCHEDULER=
false ship-dark путь из #2192) удалён — kit (_run_kit_scheduler) теперь единственный
путь. Прод уже давно на kit (USE_KIT_SCHEDULER=true), legacy был мёртвым грузом.
settings.use_kit_scheduler остаётся в конфиге (extra="ignore" защищает от startup-краха
на leftover env var), но на ветвление больше не влияет.
"""
task = asyncio.create_task(_run_kit_scheduler())
loop = asyncio.get_running_loop()
def _on_signal(signum: int) -> None:
logger.info("scheduler_main: signal %d received — requesting cooperative drain", signum)
request_shutdown()
try:
loop.add_signal_handler(signal.SIGTERM, lambda: _on_signal(signal.SIGTERM))
loop.add_signal_handler(signal.SIGINT, lambda: _on_signal(signal.SIGINT))
except NotImplementedError:
# Windows dev: signal handlers через loop не поддерживаются
logger.warning("scheduler_main: loop.add_signal_handler not supported (Windows dev)")
await _await_scheduler(task)
if shutdown_requested():
logger.info("scheduler_main: scheduler drained cleanly (SIGTERM)")
else:
logger.info("scheduler_main: scheduler task exited")
if __name__ == "__main__":
if not _should_run():
logger.warning("scheduler_main: scheduler disabled via SCHEDULER_ENABLE — exiting")
sys.exit(0)
# Best-effort FDW bootstrap: чтобы ensure_fdw_user_mapping отработал при старте
# scraper-контейнера (postgres_fdw нужен для rosreestr_dkp_import).
try:
from app.core.db import SessionLocal
from app.core.fdw import ensure_fdw_user_mapping
with SessionLocal() as db:
ensure_fdw_user_mapping(db)
except Exception:
logger.exception(
"scheduler_main: FDW user mapping bootstrap failed — cadastral queries may fail"
)
asyncio.run(_run())