diff --git a/apps/api/app/api/routes/generation_tasks.py b/apps/api/app/api/routes/generation_tasks.py old mode 100644 new mode 100755 index 36c11f709..ebedda705 --- a/apps/api/app/api/routes/generation_tasks.py +++ b/apps/api/app/api/routes/generation_tasks.py @@ -37,6 +37,43 @@ logger = logging.getLogger(__name__) router = APIRouter() +def _safe_enqueue_generation_task( + task: Any, + generation_task_repository: Any, +) -> bool: + """安全入队:send_task 失败时自动把任务标记为 failed,避免留下 pending 僵尸任务。 + + Returns: + True 表示入队成功,False 表示入队失败(已标记为 failed) + """ + try: + celery_app.send_task("worker.generate_video", args=[task.id]) + logger.info( + "[生成任务] 入队成功: task_id=%s, status=%s", + task.id, + task.status, + ) + return True + except Exception as e: + logger.error( + "[生成任务] 入队失败,标记为失败: task_id=%s error=%s", + task.id, + e, + exc_info=True, + ) + try: + task.mark_failed(f"任务入队失败: {e}") + generation_task_repository.update(task) + except Exception as update_err: + logger.error( + "[生成任务] 入队失败后更新状态也失败: task_id=%s error=%s", + task.id, + update_err, + exc_info=True, + ) + return False + + def _check_project_access(project_id: str, user_id: str, project_repository) -> None: """检查用户是否有项目访问权限""" project = project_repository.find_by_id(project_id) @@ -228,6 +265,7 @@ def create_generation_task( use_case = CreateGenerationTaskUseCase(generation_task_repository) count = request.count created_tasks = [] + failed_tasks = [] # 同批次任务共享 batch_id,用于视频查重时批次内比对 batch_id = uuid.uuid4().hex if count > 1 else "" @@ -249,19 +287,15 @@ def create_generation_task( batch_id=batch_id, ) ) - celery_app.send_task("worker.generate_video", args=[task.id]) - created_tasks.append(task) - logger.info( - "[生成任务] 入队成功: task_id=%s, status=%s, batch_id=%s", - task.id, - task.status, - batch_id, - ) + if _safe_enqueue_generation_task(task, generation_task_repository): + created_tasks.append(task) + else: + failed_tasks.append(task) except Exception as e: logger.error("[生成任务] 创建失败: %s", e, exc_info=True) raise HTTPException(status_code=500, detail="创建生成任务失败,请稍后重试或查看任务日志") - items = [_to_generation_task_response(t) for t in created_tasks] + items = [_to_generation_task_response(t) for t in created_tasks + failed_tasks] return BatchGenerationTaskResponse(items=items, total=len(items)) @@ -347,5 +381,6 @@ def retry_generation_task( asset_select_mode=getattr(task, "asset_select_mode", ""), ) ) - celery_app.send_task("worker.generate_video", args=[retried.id]) + if not _safe_enqueue_generation_task(retried, generation_task_repository): + logger.warning("[生成任务] 重试入队失败: task_id=%s", retried.id) return _to_generation_task_response(retried) diff --git a/apps/api/app/api/routes/task_center.py b/apps/api/app/api/routes/task_center.py old mode 100644 new mode 100755 index 62c97c960..13e51f480 --- a/apps/api/app/api/routes/task_center.py +++ b/apps/api/app/api/routes/task_center.py @@ -25,6 +25,35 @@ from packages.application import ( router = APIRouter() +def _safe_enqueue_generation_task( + task: Any, + generation_task_repository: Any, +) -> bool: + """安全入队:send_task 失败时自动把任务标记为 failed,避免留下 pending 僵尸任务。""" + try: + celery_app.send_task("worker.generate_video", args=[task.id]) + logger.info("[任务中心] 生成任务入队成功: task_id=%s", task.id) + return True + except Exception as e: + logger.error( + "[任务中心] 生成任务入队失败,标记为失败: task_id=%s error=%s", + task.id, + e, + exc_info=True, + ) + try: + task.mark_failed(f"任务入队失败: {e}") + generation_task_repository.update(task) + except Exception as update_err: + logger.error( + "[任务中心] 入队失败后更新状态也失败: task_id=%s error=%s", + task.id, + update_err, + exc_info=True, + ) + return False + + def _humanize_task_error(error_message: str) -> str: raw = (error_message or "").strip() if not raw: @@ -153,7 +182,8 @@ def retry_task_by_id( created_by_user_id=authenticated_user.user.id, ) ) - celery_app.send_task("worker.generate_video", args=[retried.id]) + if not _safe_enqueue_generation_task(retried, generation_task_repository): + logger.warning("[任务中心] 用户级重试入队失败: task_id=%s", retried.id) return UserTaskResponse( id=f"generation:{retried.id}", task_type="generation", @@ -235,7 +265,8 @@ def retry_project_task( created_by_user_id=authenticated_user.user.id, ) ) - celery_app.send_task("worker.generate_video", args=[retried.id]) + if not _safe_enqueue_generation_task(retried, generation_task_repository): + logger.warning("[任务中心] 项目级重试用队失败: task_id=%s", retried.id) return _generation_task_to_project_response(retried) if task_type == "ingest": job = ingest_job_repository.get(source_id)