from airshelf.celery import app from apps.ai.models import AITask from apps.ai.services import poll_video_segment from apps.projects.models import ExportJob from apps.projects.models import VideoSegment from apps.projects.services.export import run_export_job @app.task(bind=True, max_retries=120) def poll_video_segment_task(self, video_segment_id: str) -> str: segment = VideoSegment.objects.select_related("project", "project__created_by").get(id=video_segment_id) ai_task = ( segment.project.ai_tasks.filter( task_type=AITask.Type.VIDEO_SEGMENT, request_payload__video_segment_id=str(segment.id), status__in=[AITask.Status.SUBMITTED, AITask.Status.POLLING], ) .select_related("created_by") .order_by("-created_at") .first() ) if ai_task is None: return video_segment_id user = ai_task.created_by or segment.project.created_by version = poll_video_segment(video_segment=segment, user=user) if version is None and segment.status in [VideoSegment.Status.RUNNING, VideoSegment.Status.QUEUED]: raise self.retry(countdown=30) return video_segment_id @app.task(bind=True, max_retries=2) def run_export_job_task(self, export_job_id: str) -> str: try: run_export_job(export_job_id) except Exception as exc: export_job = ExportJob.objects.filter(id=export_job_id).first() if export_job: export_job.status = ExportJob.Status.FAILED export_job.error_message = str(exc) export_job.save(update_fields=["status", "error_message", "updated_at"]) raise return export_job_id QUICK_CREATE_QUEUE = "airshelf.quick" @app.task(bind=True, max_retries=0, soft_time_limit=240, time_limit=270, queue=QUICK_CREATE_QUEUE) def run_quick_script_task(self, quick_job_id: str) -> str: """脚本生成单独跑,避免把整条极速成片编排堵在一次 SSE 消费里。""" from celery.exceptions import SoftTimeLimitExceeded from apps.projects.models import QuickCreateJob from apps.projects.services.quick_create import consume_quick_script, fail_quick_create try: consume_quick_script(quick_job_id) except SoftTimeLimitExceeded: job = QuickCreateJob.objects.select_related("project").filter(id=quick_job_id).first() if job is not None and job.status not in { QuickCreateJob.Status.SUCCEEDED, QuickCreateJob.Status.FAILED, QuickCreateJob.Status.CANCELLED, }: fail_quick_create(job, "脚本生成超时,请稍后重试或进入专业模式查看") raise return quick_job_id @app.task(bind=True, max_retries=0, queue=QUICK_CREATE_QUEUE) def advance_quick_create_task(self, quick_job_id: str) -> str: """一次只推进一个可重入状态,等待型阶段通过重新入队轮询,不占 worker 睡眠。""" from apps.projects.services.quick_create import advance_quick_create next_delay = advance_quick_create(quick_job_id) if next_delay is not None: from apps.projects.services.quick_create import _claim_next_advance if not _claim_next_advance(quick_job_id, int(next_delay)): return quick_job_id try: advance_quick_create_task.apply_async( args=[quick_job_id], countdown=max(1, int(next_delay)), queue=QUICK_CREATE_QUEUE, ) except Exception as exc: # noqa: BLE001 — 重排失败不能留下永久“生成中” from apps.projects.models import QuickCreateJob from apps.projects.services.quick_create import fail_quick_create job = QuickCreateJob.objects.select_related("project").filter(id=quick_job_id).first() if job is not None: fail_quick_create(job, "生成队列暂时中断,请稍后重试", internal_error=str(exc)) raise return quick_job_id