/
asigatchov
/
vb-ai-api
Обзор
Документация
Войти
/
asigatchov
/
vb-ai-api
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
src/celery_worker.py
169 строк
5 KB
Alexander Sigatchov
update api
28 фев 2026, 18:57
28 фев 2026, 18:57
aaf14d2
Код
Авторство
О чём код?
import logging import os from pathlib import Path from celery import Celery from sqlalchemy import select from src.api.v1.rallies import service as rallies_service from src.config.settings import DB_PATH, UPLOADS_DIR from src.database.models import Project, Reel from src.database.session import init_db logger = logging.getLogger(__name__) DEFAULT_REDIS_URL = "redis://redis:6379/0" BROKER_URL = os.getenv("CELERY_BROKER_URL", os.getenv("REDIS_URL", DEFAULT_REDIS_URL)) RESULT_BACKEND = os.getenv("CELERY_RESULT_BACKEND", os.getenv("REDIS_URL", DEFAULT_REDIS_URL)) API_IMPORT_TASK_NAME = os.getenv("API_IMPORT_TASK_NAME", "api.import_rallies_from_tracks") API_SET_STATUS_TASK_NAME = os.getenv("API_SET_STATUS_TASK_NAME", "api.set_project_status") API_UPDATE_REEL_TASK_NAME = os.getenv("API_UPDATE_REEL_TASK_NAME", "api.update_reel_result") celery_app = Celery("vb_ai_api_worker", broker=BROKER_URL, backend=RESULT_BACKEND) celery_app.conf.update( task_serializer="json", accept_content=["json"], result_serializer="json", timezone="UTC", enable_utc=True, ) celery = celery_app app = celery_app _session_factory = init_db(DB_PATH) _uploads_dir = Path(os.getenv("UPLOADS_DIR", str(UPLOADS_DIR))).resolve() def _set_project_status(db, project_id: str, user_id: str, status: str) -> None: project = db.scalar(select(Project).where(Project.id == project_id, Project.user_id == user_id)) if project: project.status = status db.flush() def _set_reel_result( db, reel_id: str, project_id: str, user_id: str, status: str, url: str | None = None, title: str | None = None, rally_id: str | None = None, ) -> None: project = db.scalar(select(Project).where(Project.id == project_id, Project.user_id == user_id)) if not project: return reel = db.scalar( select(Reel) .where(Reel.id == reel_id, Reel.project_id == project_id) ) if not reel: if not rally_id: return reel = Reel( id=reel_id, project_id=project_id, rally_id=rally_id, title=title or "reel", status=status, url=url, published=False, ) db.add(reel) db.flush() return if rally_id: reel.rally_id = rally_id if title: reel.title = title reel.status = status reel.url = url db.flush() @celery_app.task(name=API_SET_STATUS_TASK_NAME) def set_project_status_task(project_id: str, user_id: str, status: str) -> dict[str, str]: db = _session_factory() try: _set_project_status(db, project_id, user_id, status) db.commit() return {"project_id": project_id, "status": status} except Exception: db.rollback() logger.exception("Failed to set project status for project=%s user=%s", project_id, user_id) raise finally: db.close() @celery_app.task(name=API_IMPORT_TASK_NAME) def import_rallies_from_tracks_task(project_id: str, user_id: str, replace_existing: bool = True) -> dict[str, int | str]: db = _session_factory() try: result = rallies_service.import_rallies_from_tracks_by_project_id( project_id=project_id, user_id=user_id, replace_existing=replace_existing, db=db, uploads_dir=_uploads_dir, ) _set_project_status(db, project_id, user_id, "finish") db.commit() logger.info( "Imported tracks for project=%s user=%s imported=%s skipped=%s", project_id, user_id, result.imported, result.skipped, ) return { "project_id": project_id, "imported": result.imported, "skipped": result.skipped, } except Exception: db.rollback() try: _set_project_status(db, project_id, user_id, "new") db.commit() except Exception: db.rollback() logger.exception("Failed to rollback project status to new for project=%s user=%s", project_id, user_id) logger.exception("Failed to import tracks for project=%s user=%s", project_id, user_id) raise finally: db.close() @celery_app.task(name=API_UPDATE_REEL_TASK_NAME) def update_reel_result_task( reel_id: str, project_id: str, user_id: str, status: str, url: str | None = None, title: str | None = None, rally_id: str | None = None, ) -> dict[str, str | None]: db = _session_factory() try: _set_reel_result(db, reel_id, project_id, user_id, status, url, title, rally_id) db.commit() return { "reel_id": reel_id, "project_id": project_id, "status": status, "url": url, "title": title, "rally_id": rally_id, } except Exception: db.rollback() logger.exception("Failed to update reel result reel=%s project=%s user=%s", reel_id, project_id, user_id) raise finally: db.close()