# Celery tasks для синхронизации автомобилей Encar. import logging from threading import Event, Thread import uuid from celery import shared_task from celery.exceptions import MaxRetriesExceededError, Retry from celery.signals import worker_process_init from redis import Redis from ..core.config import Settings from ..core.runtime_config import load_runtime_config from ..encar import EncarScraper, EncarFilters, expand_allowed_brands from ..storage.db import PersistenceService logger = logging.getLogger("encar_scraper.worker.tasks") ENCAR_SYNC_LOCK_KEY = "encar:locks:sync_listing" _persistence: PersistenceService | None = None _redis: Redis | None = None @worker_process_init.connect def _reset_globals_after_fork(**kwargs): global _persistence, _redis # После fork закрываем унаследованные TCP-сокеты родителя (иначе дети # делят одни connection pool'ы → повреждение данных). if _persistence is not None: try: _persistence.engine.dispose() except Exception: logger.debug("engine.dispose after fork failed", exc_info=True) if _redis is not None: try: _redis.close() except Exception: logger.debug("redis.close after fork failed", exc_info=True) _persistence = None _redis = None def _sync_listing_lock_ttl_seconds() -> int: settings = Settings() return max(settings.celery.task_time_limit + 120, 300) def _get_persistence() -> PersistenceService: global _persistence if _persistence is None: _persistence = PersistenceService(Settings()) return _persistence def _get_redis() -> Redis: global _redis if _redis is None: settings = Settings() _redis = Redis.from_url( settings.redis.url, decode_responses=True, socket_timeout=10, socket_connect_timeout=5, ) return _redis def _acquire_lock(redis_client: Redis, key: str, owner_token: str, ttl_seconds: int) -> bool: try: return bool(redis_client.set(key, owner_token, nx=True, ex=ttl_seconds)) except Exception as exc: logger.warning("Failed to acquire lock %s", key, exc_info=True) return False _REFRESH_LOCK_SCRIPT = """ if redis.call('GET', KEYS[1]) == ARGV[1] then return redis.call('EXPIRE', KEYS[1], tonumber(ARGV[2])) end return 0 """ _RELEASE_LOCK_SCRIPT = """ if redis.call('GET', KEYS[1]) == ARGV[1] then return redis.call('DEL', KEYS[1]) end return 0 """ def _refresh_lock_if_owner(redis_client: Redis, key: str, owner_token: str, ttl_seconds: int) -> bool | None: try: script = redis_client.register_script(_REFRESH_LOCK_SCRIPT) refreshed = script(keys=[key], args=[owner_token, int(ttl_seconds)]) return bool(refreshed) except Exception as exc: logger.warning("Failed to refresh lock %s", key, exc_info=True) return None def _release_lock_if_owner(redis_client: Redis, key: str, owner_token: str) -> None: try: script = redis_client.register_script(_RELEASE_LOCK_SCRIPT) script(keys=[key], args=[owner_token]) except Exception as exc: logger.warning("Failed to release lock %s", key, exc_info=True) def _start_lock_heartbeat( redis_client: Redis, key: str, owner_token: str, ttl_seconds: int, ) -> tuple[Event, Thread]: stop_event = Event() interval_seconds = max(5.0, min(30.0, ttl_seconds / 3)) def _heartbeat() -> None: while not stop_event.wait(interval_seconds): refreshed = _refresh_lock_if_owner(redis_client, key, owner_token, ttl_seconds) if refreshed is False: logger.warning("Lost sync_listing lock ownership for %s", owner_token) return thread = Thread(target=_heartbeat, name="encar-lock-heartbeat", daemon=True) thread.start() return stop_event, thread @shared_task( name="encar_scraper.worker.tasks.encar_sync_listing_task", bind=True, max_retries=2, default_retry_delay=60, acks_late=True, ) def encar_sync_listing_task( self, car_type: str = "all", manufacturer: str | None = None, year_from: int | None = None, year_to: int | None = None, price_max: int | None = None, lane: str = "encar", limit: int | None = None, ): """ Полная синхронизация листинга Encar. Собирает ВСЕ авто, обновляет цены/данные, помечает проданные. """ persistence = _get_persistence() persistence.create_tables() task_id = self.request.id or "unknown" owner_token = f"{task_id}:{uuid.uuid4().hex}" redis_client = _get_redis() lock_ttl = _sync_listing_lock_ttl_seconds() lock_acquired = _acquire_lock(redis_client, ENCAR_SYNC_LOCK_KEY, owner_token, lock_ttl) if not lock_acquired: logger.info("encar_sync_listing_task skipped: another sync is already running") return { "status": "skipped", "reason": "sync_already_running", "task_id": task_id, } heartbeat_stop: Event | None = None heartbeat_thread: Thread | None = None run_id: int | None = None # Планировать следующий прогон только когда текущий прогон завершён # (success или final failure после исчерпания retries). Retry сам # перепланирует задачу — не дублируем apply_async в этом случае. reschedule_next = True try: heartbeat_stop, heartbeat_thread = _start_lock_heartbeat( redis_client, ENCAR_SYNC_LOCK_KEY, owner_token, lock_ttl, ) self.update_state(state="STARTED", meta={"stage": "encar_sync_started", "task_id": task_id}) # --- Загрузка runtime_config --- rc = load_runtime_config() effective_limit = limit or rc.sync.limit effective_lane = lane or rc.sync.lane or "encar" # --- Фиксируем начало sync_run --- run_id = persistence.start_sync_run(effective_lane) car_type_map = {"all": "A", "domestic": "Y", "import": "N"} filters = EncarFilters( car_type=car_type_map.get(car_type, "A"), manufacturer=manufacturer, year_from=year_from, year_to=year_to, price_to=price_max, ) # Набор разрешённых брендов (lowercased) из runtime_config # Расширяем английские имена alias'ами (корейский + русский перевод) allowed_brands = expand_allowed_brands( set(rc.filters.brands) if rc.filters.brands else None ) excluded_brands = expand_allowed_brands( set(rc.filters.exclude_brands) if rc.filters.exclude_brands else None ) scraper = EncarScraper(persistence=persistence) result = scraper.sync_listing( limit=effective_limit, filters=filters, lane=effective_lane, redis_client=redis_client, allowed_brands=allowed_brands, excluded_brands=excluded_brands, runtime_filters=rc.filters, probe_all_photos=rc.sync.probe_all_photos, ) # --- Запись результата в sync_runs --- persistence.finish_sync_run( run_id, status="success", ids_fetched=result.get("items_collected", 0), cars_upserted=result.get("cars_synced", 0), cars_failed=result.get("cars_failed", 0), images_upserted=0, ) summary = { "task_id": task_id, "source": "encar", "total_available": result.get("total_available", 0), "items_collected": result.get("items_collected", 0), "cars_synced": result.get("cars_synced", 0), "cars_failed": result.get("cars_failed", 0), "cars_marked_sold": result.get("cars_marked_sold", 0), } logger.info( "encar_sync_listing_task completed: %d synced, %d failed, %d marked sold", summary["cars_synced"], summary["cars_failed"], summary["cars_marked_sold"], ) return {"status": "success", **summary} except Exception as exc: logger.error("encar_sync_listing_task failed: %s", exc, exc_info=True) try: if run_id is not None: persistence.finish_sync_run( run_id, status="failed", ids_fetched=0, cars_upserted=0, cars_failed=0, images_upserted=0, error_summary=str(exc)[:500], ) except Exception: logger.warning("Failed to record sync_run failure", exc_info=True) try: # Retry сам перепланирует задачу — не планируем следующий прогон reschedule_next = False raise self.retry(exc=exc) except MaxRetriesExceededError: # Retries исчерпаны — запускаем следующий прогон по расписанию reschedule_next = True raise finally: if heartbeat_stop is not None: heartbeat_stop.set() if heartbeat_thread is not None: heartbeat_thread.join(timeout=max(1.0, min(5.0, lock_ttl / 10))) if lock_acquired: _release_lock_if_owner(redis_client, ENCAR_SYNC_LOCK_KEY, owner_token) if reschedule_next: try: interval_minutes = Settings().celery.encar_beat_interval_minutes countdown = max(60, int(interval_minutes) * 60) self.apply_async( kwargs={"car_type": car_type}, countdown=countdown, queue="encar", ) logger.info( "Next encar_sync_listing_task scheduled in %d minutes", countdown // 60, ) except Exception: logger.warning("Failed to schedule next sync run", exc_info=True) @shared_task( name="encar_scraper.worker.tasks.encar_sync_vehicle_task", bind=True, max_retries=2, default_retry_delay=30, acks_late=True, ) def encar_sync_vehicle_task(self, vehicle_url: str, lane: str = "encar"): """Синхронизация одного авто Encar.""" persistence = _get_persistence() persistence.create_tables() try: scraper = EncarScraper(persistence=persistence) result = scraper.sync_vehicle(vehicle_url, lane=lane) logger.info("encar_sync_vehicle_task completed: %s", vehicle_url) return { "status": "success", "vehicle_url": vehicle_url, "origin_id": result.get("origin_id"), } except Exception as exc: logger.error("encar_sync_vehicle_task failed: %s — %s", vehicle_url, exc, exc_info=True) raise self.retry(exc=exc)