# Задачи Celery для синхронизации автомобилей и листинга IAAI. from concurrent.futures import ThreadPoolExecutor import logging from threading import Event, Thread import time import uuid from celery import shared_task from redis import Redis from ..core.config import Settings from ..scraper import IAAIScraper from ..storage.db import PersistenceService logger = logging.getLogger("iaai_scraper.worker.tasks") SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing" def _run_browser_job(func, *args, **kwargs): # Браузерный код запускаем в отдельном потоке без активного loop. with ThreadPoolExecutor(max_workers=1, thread_name_prefix="iaai-browser") as executor: future = executor.submit(func, *args, **kwargs) return future.result() def _sync_listing_lock_ttl_seconds() -> int: settings = Settings() # Небольшой запас к лимиту времени, чтобы lock снимался после сбоев. return max(settings.celery.task_time_limit + 120, 300) def _get_persistence() -> PersistenceService: return PersistenceService(Settings()) def _get_redis() -> Redis: settings = Settings() return Redis.from_url(settings.redis.url, decode_responses=True) 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 def _refresh_lock_if_owner(redis_client: Redis, key: str, owner_token: str, ttl_seconds: int) -> bool | None: try: refreshed = redis_client.eval( """ if redis.call('GET', KEYS[1]) == ARGV[1] then return redis.call('EXPIRE', KEYS[1], tonumber(ARGV[2])) end return 0 """, 1, key, 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: redis_client.eval( """ if redis.call('GET', KEYS[1]) == ARGV[1] then return redis.call('DEL', KEYS[1]) end return 0 """, 1, key, 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="sync-listing-lock-heartbeat", daemon=True) thread.start() return stop_event, thread @shared_task( name="iaai_scraper.worker.tasks.sync_vehicle_task", bind=True, max_retries=2, default_retry_delay=30, acks_late=True, ) def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"): # Скрапинг и upsert одного автомобиля. persistence = _get_persistence() persistence.create_tables() try: def _job(): with IAAIScraper() as scraper: return scraper.sync_vehicle(vehicle_url, lane=lane) result = _run_browser_job(_job) logger.info("sync_vehicle_task completed: %s", vehicle_url) return { "status": "success", "vehicle_url": vehicle_url, "db_action": result.get("db_action"), "images_upserted": result.get("images_upserted", 0), } except Exception as exc: logger.error("sync_vehicle_task failed: %s — %s", vehicle_url, exc, exc_info=True) raise self.retry(exc=exc) @shared_task( name="iaai_scraper.worker.tasks.sync_listing_task", bind=True, max_retries=1, default_retry_delay=120, acks_late=True, ) def sync_listing_task( self, make: str | None = None, model: str | None = None, lane: str = "iaai_cars", limit: int | None = None, only_new: bool | None = None, ): # Полный цикл: листинг + sync всех найденных машин. 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_acquired = False lock_ttl = _sync_listing_lock_ttl_seconds() heartbeat_stop: Event | None = None heartbeat_thread: Thread | None = None lock_acquired = _acquire_lock(redis_client, SYNC_LISTING_LOCK_KEY, owner_token, lock_ttl) if not lock_acquired: logger.info("sync_listing_task skipped: another sync is already running") return { "status": "skipped", "reason": "sync_already_running", "task_id": task_id, } try: heartbeat_stop, heartbeat_thread = _start_lock_heartbeat( redis_client, SYNC_LISTING_LOCK_KEY, owner_token, lock_ttl, ) self.update_state(state="STARTED", meta={"stage": "sync_listing_started", "task_id": task_id}) def _job(): with IAAIScraper() as scraper: return scraper.sync_listing( make=make, model=model, lane=lane, limit=limit, only_new=only_new, ) result = _run_browser_job(_job) summary = { "task_id": task_id, "run_id": result.get("run_id"), "cars_upserted": result.get("cars_upserted", 0), "cars_failed": result.get("cars_failed", 0), "images_upserted": result.get("images_upserted", 0), "skipped_existing": result.get("skipped_existing", 0), "elapsed_seconds": result.get("elapsed_seconds"), } logger.info( "sync_listing_task completed: %d upserted, %d failed", summary["cars_upserted"], summary["cars_failed"], ) return {"status": "success", **summary} except Exception as exc: logger.error("sync_listing_task failed: %s", exc, exc_info=True) raise self.retry(exc=exc) 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, SYNC_LISTING_LOCK_KEY, owner_token)