gendesign/tradein-mvp/backend/app/scheduler_main.py
bot-backend 30e3bacc5e
All checks were successful
CI Trade-In / changes (pull_request) Successful in 10s
CI / changes (pull_request) Successful in 12s
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
CI Trade-In / backend-tests (pull_request) Successful in 4m55s
fix(tradein/scheduler): SIGTERM-drain снимает с 'running' in-flight app-task'и (#3391)
Прод 07.09 02:36 UTC, первый настоящий drain после #3363: hard-cancel из
scheduler_main оборвал дрейн, и два бэкфилла (cian_detail_backfill 6167,
cian_history_backfill 6173) остались в scrape_runs со статусом 'running' —
boot-reap следующего контейнера сделал их 'zombie' (boot_reaped=true), метки
interrupted не было. interrupted=1 при дрейне писали только kit-пайплайны и
DKP-импорт: у задач, чьё тело живёт в app, ставить её было некому.

Метка ставится в единственной точке, через которую проходит любая detached
run-задача — SchedulerContext.drain_inflight: и по истечении
_CHILD_DRAIN_TIMEOUT_S, и в обработчике CancelledError (тот самый прод-путь).
run_id берётся из нового реестра {task: run_id}, который заполняет _dispatch
сразу после claim'а; claim-логика не тронута. Статус строки перечитывается
перед записью, поэтому успевший финализироваться сам прогон не
перезаписывается, а counters читаются из строки и дописываются — app-копия
mark_done их ЗАМЕНЯЕТ (#3390), голая {"interrupted": 1} стёрла бы чекпоинт.

scheduler_main: _await_scheduler возвращает признак hard-cancel'а, и строка
«scheduler drained cleanly (SIGTERM)» больше не печатается сразу за WARNING'ом
о превышении grace — на проде эти две строки стояли подряд и противоречили
друг другу.

Запас времени на запись: docker stop_grace_period 120s − _DRAIN_TIMEOUT_S 100s
= 20 с после hard-cancel'а, запись синхронная (несколько statement'ов).
2026-09-06 07:56:35 +05:00

267 lines
14 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.

"""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,
stabilize_retry_error_fingerprint,
)
def _before_send(event: dict, hint: dict) -> dict | None: # type: ignore[type-arg]
"""PR-D2: этот процесс не держит ASGI-приложения (нет `request` в event
сегодня), но payments_confirm/payments_reconcile (PR-E, тот же
`tradein-scraper` контейнер) будут звать Т-Банк API отсюда — belt-and-
suspenders на случай, если платёжные данные когда-нибудь попадут в
`request`/`extra`. Тот же обработчик на оба канала ниже — см.
app/main.py._before_send (идентичный мотив, не дублировать без причины).
PII-scrub + RetryError fingerprint-стабилизация (glitchtip-noise) идут
следом за платёжным body-wipe: этот процесс гоняет
`geocode_missing_listings` (ночной batch, сотни адресов за прогон) —
@retry-декорированные Nominatim-хелперы (app/services/geocoder.py) на
исчерпанных ретраях исторически плодили по отдельному GlitchTip issue
на КАЖДЫЙ адрес (RetryError.__str__() тащит нестабильный repr() Future).
См. sentry_scrub docstring.
"""
scrubbed = scrub_payment_request_body(event, hint) # type: ignore[arg-type]
if scrubbed is None:
return None
scrubbed = scrub_pii_event(scrubbed, hint)
if scrubbed is None:
return None
return stabilize_retry_error_fingerprint(scrubbed, hint)
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,
# #3194: процесс скрейпера — в локальных переменных кадров лежат
# прокси-креды вида user:pass, default sentry_sdk (True) приложил бы их
# к traceback открытым текстом. Так же выставлено у соседей (main.py,
# tgbot). НЕ закрывает значения в самом тексте исключения и в логах.
include_local_variables=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]) -> bool:
"""Дождаться завершения scheduler-задачи с кооперативным SIGTERM-drain'ом.
Возвращает True, если задачу пришлось хард-кансельнуть (grace истёк), False — если
она вышла сама. Вызывающий обязан различать эти исходы в логах (#3391): строка
«drained cleanly» после hard-cancel'а — ложь, а именно она печаталась на проде
07.09 сразу за WARNING'ом о превышении grace.
- Нет shutdown: scheduler_loop бесконечен → задача никогда не завершается, ждём её
как есть (процесс просто работает).
- shutdown запрошен, задача ещё бежит: bounded `wait_for(_DRAIN_TIMEOUT_S)` —
кооперативная задача докоммитит текущий unit на ближайшем checkpoint'е и выйдет
сама. Если превысила grace → hard-cancel + suppress CancelledError, чтобы
некооперативная задача не подвесила процесс за пределами docker stop_grace_period.
Пометку `interrupted` своим in-flight прогонам kit-scheduler успевает поставить в
обработчике CancelledError (SchedulerContext.mark_inflight_interrupted): между
hard-cancel'ом и SIGKILL'ом остаётся 20 с docker-grace (120s 100s).
"""
# Гонка «задача завершилась сама» против «пришёл 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 False
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")
return False
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
return True
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)")
hard_cancelled = await _await_scheduler(task)
if hard_cancelled:
# Дрейн НЕ был чистым: задача не вышла сама, её сняли. WARNING об этом уже
# напечатан в _await_scheduler — второй строкой её не «переобъявляем».
return
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())