Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 7ed0ffd5a8 |
Regular → Executable
+45
-10
@@ -37,6 +37,43 @@ logger = logging.getLogger(__name__)
|
|||||||
router = APIRouter()
|
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:
|
def _check_project_access(project_id: str, user_id: str, project_repository) -> None:
|
||||||
"""检查用户是否有项目访问权限"""
|
"""检查用户是否有项目访问权限"""
|
||||||
project = project_repository.find_by_id(project_id)
|
project = project_repository.find_by_id(project_id)
|
||||||
@@ -228,6 +265,7 @@ def create_generation_task(
|
|||||||
use_case = CreateGenerationTaskUseCase(generation_task_repository)
|
use_case = CreateGenerationTaskUseCase(generation_task_repository)
|
||||||
count = request.count
|
count = request.count
|
||||||
created_tasks = []
|
created_tasks = []
|
||||||
|
failed_tasks = []
|
||||||
# 同批次任务共享 batch_id,用于视频查重时批次内比对
|
# 同批次任务共享 batch_id,用于视频查重时批次内比对
|
||||||
batch_id = uuid.uuid4().hex if count > 1 else ""
|
batch_id = uuid.uuid4().hex if count > 1 else ""
|
||||||
|
|
||||||
@@ -249,19 +287,15 @@ def create_generation_task(
|
|||||||
batch_id=batch_id,
|
batch_id=batch_id,
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
celery_app.send_task("worker.generate_video", args=[task.id])
|
if _safe_enqueue_generation_task(task, generation_task_repository):
|
||||||
created_tasks.append(task)
|
created_tasks.append(task)
|
||||||
logger.info(
|
else:
|
||||||
"[生成任务] 入队成功: task_id=%s, status=%s, batch_id=%s",
|
failed_tasks.append(task)
|
||||||
task.id,
|
|
||||||
task.status,
|
|
||||||
batch_id,
|
|
||||||
)
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error("[生成任务] 创建失败: %s", e, exc_info=True)
|
logger.error("[生成任务] 创建失败: %s", e, exc_info=True)
|
||||||
raise HTTPException(status_code=500, detail="创建生成任务失败,请稍后重试或查看任务日志")
|
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))
|
return BatchGenerationTaskResponse(items=items, total=len(items))
|
||||||
|
|
||||||
|
|
||||||
@@ -347,5 +381,6 @@ def retry_generation_task(
|
|||||||
asset_select_mode=getattr(task, "asset_select_mode", ""),
|
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)
|
return _to_generation_task_response(retried)
|
||||||
|
|||||||
Regular → Executable
+33
-2
@@ -25,6 +25,35 @@ from packages.application import (
|
|||||||
router = APIRouter()
|
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:
|
def _humanize_task_error(error_message: str) -> str:
|
||||||
raw = (error_message or "").strip()
|
raw = (error_message or "").strip()
|
||||||
if not raw:
|
if not raw:
|
||||||
@@ -153,7 +182,8 @@ def retry_task_by_id(
|
|||||||
created_by_user_id=authenticated_user.user.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(
|
return UserTaskResponse(
|
||||||
id=f"generation:{retried.id}",
|
id=f"generation:{retried.id}",
|
||||||
task_type="generation",
|
task_type="generation",
|
||||||
@@ -235,7 +265,8 @@ def retry_project_task(
|
|||||||
created_by_user_id=authenticated_user.user.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 _generation_task_to_project_response(retried)
|
return _generation_task_to_project_response(retried)
|
||||||
if task_type == "ingest":
|
if task_type == "ingest":
|
||||||
job = ingest_job_repository.get(source_id)
|
job = ingest_job_repository.get(source_id)
|
||||||
|
|||||||
Reference in New Issue
Block a user