"""Celery tasks для IAAI.""" import json import logging from datetime import datetime, timezone from celery import shared_task from ..core.config import Settings from ..scraper import IAAIScraper from ..storage.db import PersistenceService logger = logging.getLogger("iaai_scraper.worker.tasks") def _get_persistence() -> PersistenceService: return PersistenceService(Settings()) @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 одного автомобиля.""" task_id = self.request.id persistence = _get_persistence() persistence.create_tables() # Регистрируем запуск. persistence.update_scrape_task( task_id, status="running", started_at=datetime.now(timezone.utc), ) try: with IAAIScraper() as scraper: result = scraper.sync_vehicle(vehicle_url, lane=lane) persistence.update_scrape_task( task_id, status="success", finished_at=datetime.now(timezone.utc), result_summary=json.dumps({ "car_upserted": True, "db_action": result.get("db_action"), "images_upserted": result.get("images_upserted", 0), "elapsed_seconds": result.get("elapsed_seconds"), }, default=str), ) 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: persistence.update_scrape_task( task_id, status="failed", finished_at=datetime.now(timezone.utc), error_message=str(exc)[:2000], ) logger.error("sync_vehicle_task failed: %s — %s", vehicle_url, exc) 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 всех найденных машин.""" task_id = self.request.id persistence = _get_persistence() persistence.create_tables() persistence.update_scrape_task( task_id, status="running", started_at=datetime.now(timezone.utc), ) try: with IAAIScraper() as scraper: result = scraper.sync_listing( make=make, model=model, lane=lane, limit=limit, only_new=only_new, ) summary = { "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"), } persistence.update_scrape_task( task_id, status="success", finished_at=datetime.now(timezone.utc), result_summary=json.dumps(summary, default=str), ) logger.info( "sync_listing_task completed: %d upserted, %d failed", summary["cars_upserted"], summary["cars_failed"], ) return {"status": "success", **summary} except Exception as exc: persistence.update_scrape_task( task_id, status="failed", finished_at=datetime.now(timezone.utc), error_message=str(exc)[:2000], ) logger.error("sync_listing_task failed: %s", exc) raise self.retry(exc=exc)