from __future__ import annotations from app.models.base import async_session from app.services.hot_opening_replicate_service import run_image_prompt_optimize, run_video_prompt_optimize from app.tasks.async_runner import run_async from app.tasks.celery_app import celery_app async def _run_image_prompt(project_id: str, step_id: str | None = None): async with async_session() as db: await run_image_prompt_optimize(db, project_id=project_id, step_id=step_id) await db.commit() async def _run_video_prompt(project_id: str, step_id: str | None = None): async with async_session() as db: await run_video_prompt_optimize(db, project_id=project_id, step_id=step_id) await db.commit() if celery_app: @celery_app.task(name="hot_opening.start_image_prompt_optimize", bind=True, max_retries=3, default_retry_delay=30) def start_image_prompt_optimize(self, project_id: str, step_id: str | None = None): """手动触发后的图片 AI 提词任务。 该任务路由到现有 gen_chatapi_create 队列,不需要新增 hot_opening worker。 """ return run_async(_run_image_prompt(project_id, step_id)) @celery_app.task(name="hot_opening.start_video_prompt_optimize", bind=True, max_retries=3, default_retry_delay=30) def start_video_prompt_optimize(self, project_id: str, step_id: str | None = None): """手动触发后的视频 AI 提词任务。 该任务路由到现有 gen_chatapi_create 队列,不需要新增 hot_opening worker。 """ return run_async(_run_video_prompt(project_id, step_id)) else: class _DisabledTask: def delay(self, *args, **kwargs): raise RuntimeError("Celery is disabled") def apply_async(self, *args, **kwargs): raise RuntimeError("Celery is disabled") start_image_prompt_optimize = _DisabledTask() start_video_prompt_optimize = _DisabledTask()