114 lines
5.1 KiB
Python
114 lines
5.1 KiB
Python
"""Celery app — single source of truth для Celery configuration.
|
||
|
||
Beat schedule build → app/workers/beat_schedule.py.
|
||
Worker lifecycle hooks (process_init, worker_ready) → app/workers/lifecycle.py.
|
||
"""
|
||
|
||
import logging
|
||
import os
|
||
|
||
import sentry_sdk
|
||
from celery import Celery
|
||
from sentry_sdk.integrations.celery import CeleryIntegration
|
||
from sentry_sdk.integrations.httpx import HttpxIntegration
|
||
from sentry_sdk.integrations.logging import LoggingIntegration
|
||
from sentry_sdk.integrations.sqlalchemy import SqlalchemyIntegration
|
||
|
||
from app.core.config import settings
|
||
from app.observability.sentry_scrub import scrub_event
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# SDK инициализируется в обоих процессах (FastAPI-сервер и Celery-воркер),
|
||
# чтобы события из тасков попадали в GlitchTip. SDK безопасен для двойного
|
||
# вызова — повторный sentry_sdk.init() в одном процессе заменяет клиента.
|
||
if settings.glitchtip_dsn:
|
||
# before_send И before_send_transaction — ОБА на scrub_event (#2457-review,
|
||
# см. app/main.py и sentry_scrub.py module docstring): до этого фикса worker
|
||
# вообще не скрабил error-события (тут before_send не было), а
|
||
# before_send_transaction был на голом scrub_sensitive_query (только URL) —
|
||
# оба канала пропускали PII.
|
||
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=settings.glitchtip_traces_sample_rate,
|
||
profiles_sample_rate=0.0,
|
||
send_default_pii=False,
|
||
# Локальные переменные кадров стека НЕ уходят в мониторинг (#2753) — см.
|
||
# app/main.py: скраб сверяет ИМЕНА ключей, а имя переменной произвольно.
|
||
# В воркере вектор шире: задачи держат в кадрах сырые ответы источников.
|
||
include_local_variables=False,
|
||
before_send=scrub_event,
|
||
before_send_transaction=scrub_event,
|
||
integrations=[
|
||
CeleryIntegration(monitor_beat_tasks=True),
|
||
SqlalchemyIntegration(),
|
||
HttpxIntegration(),
|
||
LoggingIntegration(level=logging.INFO, event_level=logging.ERROR),
|
||
],
|
||
)
|
||
logger.info(
|
||
"GlitchTip SDK initialised in Celery worker (env=%s)",
|
||
settings.environment,
|
||
)
|
||
|
||
celery_app = Celery(
|
||
"gendesign",
|
||
broker=settings.redis_url,
|
||
backend=settings.redis_url,
|
||
include=[
|
||
"app.workers.tasks.scrape_kn",
|
||
"app.workers.tasks.scrape_kn_catalog_objects",
|
||
"app.workers.tasks.scrape_kn_catalog_flats",
|
||
"app.workers.tasks.refresh_analytics",
|
||
"app.workers.tasks.scrape_objective",
|
||
"app.workers.tasks.objective_etl",
|
||
"app.workers.tasks.nspd_geo",
|
||
"app.workers.tasks.nspd_sync",
|
||
"app.workers.tasks.poi_sync",
|
||
"app.workers.tasks.noise_sync",
|
||
"app.workers.tasks.utility_infrastructure_sync",
|
||
"app.workers.tasks.pzz_sync",
|
||
"app.workers.tasks.scrape_cadastre",
|
||
"app.workers.tasks.ekburg_permits_sync",
|
||
"app.workers.tasks.cbr_macro_sync",
|
||
"app.workers.tasks.rosstat_macro_sync",
|
||
"app.workers.tasks.refresh_quarter_price_index",
|
||
"app.workers.tasks.etl_newbuilding_crossload",
|
||
"app.workers.tasks.supply_layers_refresh",
|
||
"app.workers.tasks.location_refresh",
|
||
"app.workers.tasks.forecast",
|
||
"app.workers.tasks.full_report",
|
||
"app.workers.tasks.ird_harvest",
|
||
"app.workers.tasks.ekb_krt_sync",
|
||
"app.workers.tasks.gknspecial_harvest",
|
||
"app.workers.tasks.opportunity_harvest",
|
||
"app.workers.tasks.planning_harvest",
|
||
"app.workers.tasks.zone_regulation_refresh",
|
||
"app.workers.tasks.backfill_zone_regulations",
|
||
"app.workers.tasks.reservation_ingest",
|
||
"app.workers.tasks.genplan_zones_sync",
|
||
"app.workers.tasks.ekb_ppt_tep_sync",
|
||
"app.workers.tasks.krt_geometry_sync",
|
||
"app.workers.tasks.okn_objects_sync",
|
||
"app.workers.tasks.pat_subzones_load",
|
||
"app.workers.tasks.izyatie_ocr_ingest",
|
||
"app.workers.tasks.developer_registry_refresh",
|
||
"app.workers.tasks.refresh_layout_velocity",
|
||
"app.workers.tasks.riasurt_sverdl_harvest",
|
||
"app.workers.tasks.mv_sales_tracker_refresh",
|
||
"app.workers.tasks.scrape_freshness_check",
|
||
"app.workers.tasks.connection_capacity_sync",
|
||
"app.workers.tasks.gisogd_permits_sync",
|
||
],
|
||
)
|
||
celery_app.conf.timezone = "Europe/Moscow"
|
||
|
||
# Apply beat schedule (DB → fallback → hardcoded entries)
|
||
from app.workers.beat_schedule import build_beat_schedule # noqa: E402
|
||
|
||
celery_app.conf.beat_schedule = build_beat_schedule()
|
||
|
||
# Register lifecycle hooks (import for side-effect — signal decorator registration)
|
||
from app.workers import lifecycle # noqa: E402, F401
|