# Celery-приложение и периодические задачи. import json import logging import os import time from celery import Celery from celery.signals import worker_process_init, worker_ready, setup_logging as celery_setup_logging from redis import Redis from ..core.config import settings from ..core.logs import setup_logging from .constants import ( GLOBAL_DB_PROGRESS_TS_KEY, GLOBAL_PROGRESS_TS_KEY, MOBILEDE_RUNTIME_SEGMENTS_TASK, MOBILEDE_SYNC_QUEUE, MOBILEDE_SYNC_TASK_NAME, ) logger = logging.getLogger("mobilede_scraper.worker.celery_app") STARTUP_SYNC_DISPATCH_KEY = "mobilede:state:startup_sync_dispatched" PROGRESS_KEY_PREFIX = "mobilede:state:task_progress:" MOBILEDE_SYNC_DETAIL_TASK = "mobilede.sync_detail" MOBILEDE_ENRICH_IMAGES_TASK = "mobilede.enrich_images_batch" STARTUP_SYNC_DISPATCH_TTL_SECONDS = 10 * 60 STARTUP_PROGRESS_MAX_AGE_SECONDS = 180 STARTUP_GLOBAL_PROGRESS_MAX_AGE_SECONDS = 300 def _env_bool(name: str, default: bool) -> bool: raw = os.getenv(name, "true" if default else "false").strip().lower() return raw in {"1", "true", "yes", "on"} MOBILEDE_BEAT_SYNC_ENABLED = _env_bool("MOBILEDE_BEAT_SYNC_ENABLED", True) def _runtime_sync_kwargs() -> dict[str, bool | float]: return { "delay_seconds": float(os.getenv("MOBILEDE_REQUEST_DELAY_SECONDS", "0.7")), "use_cursor": _env_bool("MOBILEDE_CURSOR_ENABLED", True), "continuous": _env_bool("MOBILEDE_CONTINUOUS_SYNC_ENABLED", True), } def _runtime_sync_expires_seconds() -> float: return settings.celery.beat_sync_interval_minutes * 60.0 def _has_fresh_active_progress( redis_client: Redis, *, max_age_seconds: int = STARTUP_PROGRESS_MAX_AGE_SECONDS, ) -> bool: now = int(time.time()) try: for raw_key in redis_client.scan_iter(f"{PROGRESS_KEY_PREFIX}*"): payload = redis_client.get(raw_key) if not payload: continue try: progress = json.loads(payload) except (TypeError, ValueError): continue ts = int(progress.get("ts") or 0) stage = str(progress.get("stage") or "") if ts > 0 and now - ts <= max_age_seconds and stage not in {"segment_done", "sync_done", "failed"}: return True except Exception: logger.warning("Failed to inspect startup progress keys", exc_info=True) return False def _has_recent_global_progress( redis_client: Redis, *, max_age_seconds: int = STARTUP_GLOBAL_PROGRESS_MAX_AGE_SECONDS, ) -> bool: now = int(time.time()) try: progress_ts = int(redis_client.get(GLOBAL_PROGRESS_TS_KEY) or 0) db_progress_ts = int(redis_client.get(GLOBAL_DB_PROGRESS_TS_KEY) or 0) except Exception: logger.warning("Failed to inspect global startup progress keys", exc_info=True) return False freshest_ts = max(progress_ts, db_progress_ts) return freshest_ts > 0 and now - freshest_ts <= max_age_seconds def _has_live_startup_progress(redis_client: Redis, *, queue_len: int) -> bool: if _has_recent_global_progress(redis_client): return True if queue_len <= 0: return False return _has_fresh_active_progress(redis_client) def _claim_startup_dispatch(redis_client: Redis, *, reset_stale: bool) -> bool: if redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=STARTUP_SYNC_DISPATCH_TTL_SECONDS): return True if not reset_stale: return False redis_client.delete(STARTUP_SYNC_DISPATCH_KEY) return bool( redis_client.set( STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=STARTUP_SYNC_DISPATCH_TTL_SECONDS, ) ) @celery_setup_logging.connect def _configure_logging(loglevel=None, **kwargs): # Пишем логи Celery в stderr. level = settings.log_level if settings.log_level else "INFO" setup_logging(level, None) @worker_process_init.connect def _on_worker_process_init(**kwargs): # Повторно настраиваем логирование после fork. level = settings.log_level if settings.log_level else "INFO" setup_logging(level, None) def _broker_url() -> str: return settings.celery.broker_url or settings.redis.url def _result_backend() -> str: return settings.celery.result_backend or settings.redis.url celery_app = Celery( "mobilede_scraper", broker=_broker_url(), backend=_result_backend(), ) # Если hard limit слишком большой, сжимаем его до soft + 120. _soft = settings.celery.task_soft_time_limit _hard = settings.celery.task_time_limit _max_hard = _soft + 120 if _soft else _hard if _hard > _max_hard: logger.info( "CELERY_TASK_TIME_LIMIT=%d too far from CELERY_TASK_SOFT_TIME_LIMIT=%d; " "clamping hard limit to %d", _hard, _soft, _max_hard, ) _hard = _max_hard beat_schedule = {} if MOBILEDE_BEAT_SYNC_ENABLED: beat_schedule = { "periodic-mobilede-sync-search": { "task": MOBILEDE_RUNTIME_SEGMENTS_TASK, "schedule": _runtime_sync_expires_seconds(), "args": (), "kwargs": _runtime_sync_kwargs(), "options": { "queue": MOBILEDE_SYNC_QUEUE, "expires": _runtime_sync_expires_seconds(), }, } } celery_app.conf.update( task_serializer="json", accept_content=["json"], result_serializer="json", timezone="UTC", enable_utc=True, task_soft_time_limit=_soft, task_time_limit=_hard, task_acks_late=True, task_reject_on_worker_lost=True, task_track_started=True, worker_concurrency=settings.celery.worker_concurrency, worker_max_tasks_per_child=settings.celery.worker_max_tasks_per_child, worker_pool=settings.celery.worker_pool, worker_prefetch_multiplier=1, broker_connection_retry_on_startup=True, broker_transport_options={ "visibility_timeout": settings.celery.broker_visibility_timeout, }, result_expires=86400, worker_redirect_stdouts=False, worker_hijack_root_logger=False, beat_schedule=beat_schedule, task_routes={ MOBILEDE_RUNTIME_SEGMENTS_TASK: {"queue": MOBILEDE_SYNC_QUEUE}, MOBILEDE_SYNC_TASK_NAME: {"queue": MOBILEDE_SYNC_QUEUE}, MOBILEDE_SYNC_DETAIL_TASK: {"queue": MOBILEDE_SYNC_QUEUE}, MOBILEDE_ENRICH_IMAGES_TASK: {"queue": MOBILEDE_SYNC_QUEUE}, "mobilede_scraper.worker.tasks.*": {"queue": MOBILEDE_SYNC_QUEUE}, }, ) celery_app.autodiscover_tasks(["mobilede_scraper.worker"]) @worker_ready.connect def _on_worker_ready(**kwargs): """При старте worker отправляем первый canonical runtime sync, если очередь пуста.""" if not _env_bool("MOBILEDE_STARTUP_SYNC_ENABLED", True): logger.info("Worker ready: startup sync dispatch disabled by MOBILEDE_STARTUP_SYNC_ENABLED") return should_dispatch = False redis_client = None try: redis_client = Redis.from_url( settings.redis.url, decode_responses=True, socket_connect_timeout=settings.redis.socket_connect_timeout_seconds, socket_timeout=settings.redis.socket_timeout_seconds, health_check_interval=settings.redis.health_check_interval_seconds, retry_on_timeout=True, ) queue_len = int(redis_client.llen(MOBILEDE_SYNC_QUEUE) or 0) has_live_progress = _has_live_startup_progress(redis_client, queue_len=queue_len) if queue_len > 0 and has_live_progress: logger.info("Worker ready: MOBILEDE_sync queue already has %d task(s); skip startup dispatch", queue_len) return has_fresh_progress = _has_fresh_active_progress(redis_client) if queue_len > 0 and not has_live_progress: logger.warning( "Worker ready: MOBILEDE_sync queue has %d task(s), but no fresh progress is visible; forcing runtime sync dispatch", queue_len, ) elif queue_len <= 0 and has_fresh_progress: logger.info("Worker ready: queue is empty but active progress is still visible; relying on startup dedupe") should_dispatch = _claim_startup_dispatch( redis_client, reset_stale=not has_live_progress, ) if should_dispatch and not has_live_progress: logger.info("Worker ready: claimed startup dispatch after stale or missing dedupe state") except Exception: logger.warning("Worker ready startup sync dedupe check failed; skipping immediate dispatch", exc_info=True) return finally: if redis_client is not None: try: redis_client.close() except Exception: pass if not should_dispatch: logger.info("Worker ready immediate sync already dispatched recently; skipping duplicate enqueue") return logger.info("Worker ready - dispatching initial mobile.de sync_runtime_segments task") celery_app.send_task( MOBILEDE_RUNTIME_SEGMENTS_TASK, kwargs=_runtime_sync_kwargs(), queue=MOBILEDE_SYNC_QUEUE, expires=_runtime_sync_expires_seconds(), )