gendesign/backend/app/workers/celery_app.py
bot-backend d7ccf48000
All checks were successful
CI / changes (pull_request) Successful in 9s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI Trade-In / changes (pull_request) Successful in 9s
CI Trade-In / backend-tests (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Successful in 2m49s
CI / backend-tests (pull_request) Successful in 15m6s
fix(ptica): скраб ПДн перед отправкой в мониторинг + честная подпись НДС в отчётах (#2457)
Портирован PII-scrub механизм МЕРЫ (scrub_pii_event, ключи client_name/
phone/email/name) в backend/app/observability/sentry_scrub.py и подключен
как before_send в app/main.py и app/workers/celery_app.py — раньше worker
вообще не скрабил error-события, только transaction-spans (URL-secrets).
send_default_pii=False эти поля не закрывает (проверено на sentry-sdk 2.58).

full_report_docx.py / full_report_html.py: "НДС (паркинг)" -> "НДС (паркинг
+ коммерция)" — vat_rub системно включает office_value_added (коммерция),
подпись занижала состав суммы (#2457).
2026-08-06 20:47:53 +03:00

116 lines
5 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_pii_event, scrub_sensitive_query
logger = logging.getLogger(__name__)
# SDK инициализируется в обоих процессах (FastAPI-сервер и Celery-воркер),
# чтобы события из тасков попадали в GlitchTip. SDK безопасен для двойного
# вызова — повторный sentry_sdk.init() в одном процессе заменяет клиента.
if settings.glitchtip_dsn:
def _before_send(event: dict[str, object], hint: dict[str, object]) -> dict[str, object] | None:
"""Композиция PII-scrub (client_name/phone/email/name из request.data/
extra/contexts, аудит-фикс) + URL query-string secrets redact — см.
app/main.py._before_send (идентичная композиция; до этого фикса worker
вообще не скрабил error-события, только transaction-spans)."""
scrubbed = scrub_pii_event(event, hint) # type: ignore[arg-type]
if scrubbed is None:
return None
return scrub_sensitive_query(scrubbed, hint) # type: ignore[arg-type,return-value]
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,
before_send=_before_send,
before_send_transaction=scrub_sensitive_query,
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