507ca4d6b8
CI/CD Pipeline / Production Browser E2E (push) Failing after 1515h7m36s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 1515h7m36s
CI/CD Pipeline / Deploy Production (push) Failing after 1515h7m37s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Failing after 1515h7m37s
CI/CD Pipeline / Build Production Web Image (push) Failing after 1515h7m39s
CI/CD Pipeline / Build Staging Worker Image (push) Failing after 1515h7m39s
CI/CD Pipeline / Build Production Worker Image (push) Failing after 1515h7m39s
CI/CD Pipeline / Build Staging API Image (push) Failing after 1515h7m39s
CI/CD Pipeline / Build Production API Image (push) Failing after 1515h7m39s
CI/CD Pipeline / Validate Code Quality And Tests (push) Has been skipped
CI/CD Pipeline / Unit Tests (push) Has been skipped
CI/CD Pipeline / Integration Tests (push) Has been skipped
CI/CD Pipeline / Frontend Lint (push) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (push) Failing after 1515h39m18s
CI/CD Pipeline / Build Staging Web Image (push) Failing after 1515h39m21s
P0修复: - 触发生成时将plan.config.asset_ids传递到GenerationTask - 同时补全project_id字段(之前为空字符串) P1修复: - result_count: 剪辑计划多片段合成1个成片,改为1(之前误存为clip数) - total_duration: 渲染成功后回写实际输出时长到EditPlan.total_duration - logs: 新增5个关键节点日志埋点(render_start/download_done/download_failed/render_complete/render_failed) - 下载完成后进度推进到30%,便于前端展示素材下载中状态
400 lines
15 KiB
Python
Executable File
400 lines
15 KiB
Python
Executable File
"""剪辑计划生成相关 API 端点。
|
|
|
|
从 edit_plans.py 拆分,包含:
|
|
- POST /{plan_id}/generate 触发剪辑渲染生成
|
|
- GET /{plan_id}/generation-status 查询生成进度
|
|
- GET /{plan_id}/generations 查询关联的生成记录
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from typing import Any
|
|
|
|
from app.api.routes._helpers import check_project_access
|
|
from app.api.routes.edit_plans import (
|
|
ClipStatusItem,
|
|
EditPlanGenerateResponse,
|
|
EditPlanGenerationsResponse,
|
|
EditPlanGenerationStatusResponse,
|
|
)
|
|
from app.auth import AuthenticatedUser, get_current_user
|
|
from app.core.celery_app import celery_app
|
|
from app.core.task_enqueue import GLOBAL_PENDING_LIMIT, USER_PENDING_LIMIT
|
|
from app.dependencies import get_asset_library_repository, get_asset_repository, get_db_session, get_project_repository
|
|
from app.services import EditPlanService
|
|
from fastapi import APIRouter, Depends, HTTPException, status
|
|
from sqlalchemy.orm import Session
|
|
|
|
from packages.adapters.sqlalchemy_impl.generation_task_repository import (
|
|
SQLAlchemyGenerationTaskRepository,
|
|
)
|
|
from packages.adapters.sqlalchemy_impl.template_clip_config_repository import (
|
|
SQLAlchemyTemplateClipConfigRepository,
|
|
)
|
|
from packages.adapters.sqlalchemy_impl.template_repository import (
|
|
SQLAlchemyTemplateRepository,
|
|
)
|
|
from packages.application.generation_tasks import (
|
|
CreateGenerationTaskCommand,
|
|
CreateGenerationTaskUseCase,
|
|
)
|
|
from packages.domain.edit_plan import EditPlanStatus
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
router = APIRouter()
|
|
|
|
|
|
def _auto_fallback_draft_to_editing(svc: EditPlanService, plan_id: str, plan_check) -> None:
|
|
"""自动兜底 1: draft → editing"""
|
|
if plan_check.status == EditPlanStatus.DRAFT:
|
|
logger.info("自动兜底: plan=%s draft→editing", plan_id)
|
|
svc.transition_status(plan_id, EditPlanStatus.EDITING)
|
|
|
|
|
|
def _auto_fallback_copy_template_clips(svc: EditPlanService, plan_id: str, plan_check, db: Session) -> None:
|
|
"""自动兜底 2: 无片段 + 有 template_id → 从模板复制片段配置"""
|
|
existing_clips = svc.count_clips(plan_id)
|
|
if existing_clips == 0 and plan_check.template_id:
|
|
logger.info(
|
|
"自动兜底: plan=%s 无片段,从模板 %s 复制片段配置",
|
|
plan_id,
|
|
plan_check.template_id,
|
|
)
|
|
clip_config_repo = SQLAlchemyTemplateClipConfigRepository(db)
|
|
configs = clip_config_repo.list_by_template(plan_check.template_id)
|
|
if configs:
|
|
for cfg in configs:
|
|
svc.create_clip(
|
|
plan_id=plan_id,
|
|
clip_type=cfg.clip_type.value if hasattr(cfg.clip_type, "value") else cfg.clip_type,
|
|
order=cfg.order,
|
|
template_clip_config_id=cfg.id,
|
|
duration=cfg.default_duration,
|
|
transition_effect=(
|
|
cfg.transition_effect.value
|
|
if hasattr(cfg.transition_effect, "value")
|
|
else cfg.transition_effect
|
|
),
|
|
)
|
|
logger.info("自动兜底: plan=%s 从新模型 template_clip_configs 复制了 %d 个片段", plan_id, len(configs))
|
|
else:
|
|
tpl_repo = SQLAlchemyTemplateRepository(db)
|
|
segments = tpl_repo.list_segments(plan_check.template_id)
|
|
for seg in segments:
|
|
avg_duration = (seg.duration_min + seg.duration_max) / 2
|
|
svc.create_clip(
|
|
plan_id=plan_id,
|
|
clip_type="main",
|
|
order=seg.segment_order,
|
|
duration=avg_duration,
|
|
config={
|
|
"material_type": seg.material_type or "",
|
|
"template_segment_id": seg.id,
|
|
},
|
|
)
|
|
logger.info("自动兜底: plan=%s 从旧模型 template_segments 复制了 %d 个片段", plan_id, len(segments))
|
|
|
|
|
|
def _auto_fallback_assign_assets(
|
|
svc: EditPlanService,
|
|
plan_id: str,
|
|
plan_check,
|
|
) -> list:
|
|
"""自动兜底 3: 为没有素材的片段分配素材。返回剩余无素材片段列表。"""
|
|
all_clips = svc.list_clips(plan_id)
|
|
clips_without_asset = [c for c in all_clips if not c.asset_id]
|
|
config_asset_ids = (plan_check.config or {}).get("asset_ids", [])
|
|
|
|
if clips_without_asset and config_asset_ids:
|
|
logger.info(
|
|
"自动兜底3: plan=%s 为 %d 个无素材片段分配 %d 个指定素材",
|
|
plan_id,
|
|
len(clips_without_asset),
|
|
len(config_asset_ids),
|
|
)
|
|
for i, clip in enumerate(clips_without_asset):
|
|
asset_idx = i % len(config_asset_ids)
|
|
svc.assign_asset(clip.id, config_asset_ids[asset_idx])
|
|
logger.info("自动兜底3: plan=%s 素材分配完成", plan_id)
|
|
clips_without_asset = []
|
|
|
|
return clips_without_asset
|
|
|
|
|
|
def _auto_fallback_auto_material_mode(
|
|
svc: EditPlanService,
|
|
plan_id: str,
|
|
plan_check,
|
|
clips_without_asset: list,
|
|
asset_library_repo: Any,
|
|
asset_repo: Any,
|
|
) -> None:
|
|
"""自动兜底 4: 自动素材模式 → 从项目默认视频素材库选取"""
|
|
if not clips_without_asset:
|
|
return
|
|
material_mode = (plan_check.config or {}).get("material_mode", "manual")
|
|
if material_mode != "auto" or not plan_check.project_id:
|
|
return
|
|
|
|
import random
|
|
|
|
logger.info(
|
|
"自动兜底4: plan=%s 自动素材模式,从项目素材库选取素材 (%d 个片段需要)",
|
|
plan_id,
|
|
len(clips_without_asset),
|
|
)
|
|
libs = asset_library_repo.find_by_project(plan_check.project_id)
|
|
video_lib = None
|
|
for lib in libs:
|
|
lib_kind = lib.kind.value if hasattr(lib.kind, "value") else lib.kind
|
|
if lib_kind == "video":
|
|
video_lib = lib
|
|
break
|
|
|
|
if video_lib:
|
|
assets = asset_repo.find_by_library(video_lib.id)
|
|
ready_videos = [
|
|
a
|
|
for a in assets
|
|
if (a.status.value if hasattr(a.status, "value") else a.status) == "ready"
|
|
and a.mime_type
|
|
and a.mime_type.startswith("video")
|
|
]
|
|
if ready_videos:
|
|
random.shuffle(ready_videos)
|
|
for i, clip in enumerate(clips_without_asset):
|
|
asset = ready_videos[i % len(ready_videos)]
|
|
svc.assign_asset(clip.id, asset.id)
|
|
logger.info(
|
|
"自动兜底4: plan=%s 从素材库 %s 分配了 %d 个素材给 %d 个片段",
|
|
plan_id,
|
|
video_lib.name,
|
|
len(ready_videos),
|
|
len(clips_without_asset),
|
|
)
|
|
else:
|
|
logger.warning("自动兜底4: plan=%s 素材库无可用视频素材", plan_id)
|
|
else:
|
|
logger.warning("自动兜底4: plan=%s 项目无视频素材库", plan_id)
|
|
|
|
|
|
def _check_queue_limits(gen_task_repo, user_id: str) -> None:
|
|
"""队列限流预检查"""
|
|
try:
|
|
has_count = hasattr(gen_task_repo, "count_pending_by_user") and hasattr(gen_task_repo, "count_pending_total")
|
|
if has_count:
|
|
user_pending = gen_task_repo.count_pending_by_user(user_id)
|
|
global_pending = gen_task_repo.count_pending_total()
|
|
if user_pending >= USER_PENDING_LIMIT:
|
|
raise HTTPException(
|
|
status_code=429,
|
|
detail=f"您的待处理任务过多(当前 {user_pending}/{USER_PENDING_LIMIT}),请等待完成后再提交",
|
|
)
|
|
if global_pending >= GLOBAL_PENDING_LIMIT:
|
|
raise HTTPException(
|
|
status_code=503,
|
|
detail="系统繁忙,请稍后再试",
|
|
)
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.warning("[队列限流] 剪辑计划限流检查失败,跳过: %s", e)
|
|
|
|
|
|
@router.post("/{plan_id}/generate", response_model=EditPlanGenerateResponse)
|
|
def generate_plan(
|
|
plan_id: str,
|
|
db: Session = Depends(get_db_session),
|
|
current_user: AuthenticatedUser = Depends(get_current_user),
|
|
project_repository: Any = Depends(get_project_repository),
|
|
asset_library_repo: Any = Depends(get_asset_library_repository),
|
|
asset_repo: Any = Depends(get_asset_repository),
|
|
) -> EditPlanGenerateResponse:
|
|
"""触发剪辑计划渲染生成
|
|
|
|
前置条件:计划状态必须为 editing,且至少有一个片段。
|
|
流程:
|
|
1. 验证计划状态为 editing
|
|
2. 将 pending 片段标记为 ready
|
|
3. 创建 GenerationTask
|
|
4. 调度 Celery 任务 worker.render_edit_plan
|
|
5. 将计划状态流转为 rendering
|
|
"""
|
|
svc = EditPlanService(db)
|
|
plan_check = svc.get_plan(plan_id)
|
|
if plan_check is None:
|
|
raise HTTPException(status_code=404, detail=f"剪辑计划不存在: {plan_id}")
|
|
if plan_check.project_id:
|
|
check_project_access(plan_check.project_id, current_user.user.id, project_repository)
|
|
|
|
# 自动兜底流程
|
|
_auto_fallback_draft_to_editing(svc, plan_id, plan_check)
|
|
_auto_fallback_copy_template_clips(svc, plan_id, plan_check, db)
|
|
clips_without_asset = _auto_fallback_assign_assets(svc, plan_id, plan_check)
|
|
_auto_fallback_auto_material_mode(svc, plan_id, plan_check, clips_without_asset, asset_library_repo, asset_repo)
|
|
|
|
# 检查是否可生成
|
|
try:
|
|
can_gen, reason = svc.can_generate(plan_id)
|
|
except ValueError as exc:
|
|
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=str(exc)) from exc
|
|
if not can_gen:
|
|
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=reason)
|
|
|
|
# 核心生成流程
|
|
try:
|
|
clip_count = svc.mark_clips_ready(plan_id)
|
|
|
|
gen_task_repo = SQLAlchemyGenerationTaskRepository(db)
|
|
user_id = current_user.user.id
|
|
_check_queue_limits(gen_task_repo, user_id)
|
|
|
|
gen_task_use_case = CreateGenerationTaskUseCase(gen_task_repo)
|
|
plan = svc.get_plan_or_raise(plan_id)
|
|
# 从 plan.config 中读取 asset_ids 并传递给 GenerationTask
|
|
config_asset_ids = (plan.config or {}).get("asset_ids", [])
|
|
gen_task = gen_task_use_case.execute(
|
|
CreateGenerationTaskCommand(
|
|
project_id=plan.project_id or "",
|
|
template_id=plan.template_id,
|
|
created_by_user_id=current_user.user.id,
|
|
source_edit_plan_id=plan_id,
|
|
asset_ids=list(config_asset_ids) if config_asset_ids else [],
|
|
)
|
|
)
|
|
|
|
svc.update_plan_config(plan_id, {"generation_task_id": gen_task.id})
|
|
svc.transition_status(plan_id, EditPlanStatus.RENDERING)
|
|
celery_app.send_task("worker.render_edit_plan", args=[plan_id])
|
|
|
|
updated_plan = svc.get_plan_or_raise(plan_id)
|
|
|
|
logger.info(
|
|
"触发剪辑计划生成: plan_id=%s gen_task_id=%s clips=%d by user=%s",
|
|
plan_id,
|
|
gen_task.id,
|
|
clip_count,
|
|
current_user.user.id,
|
|
)
|
|
|
|
return EditPlanGenerateResponse(
|
|
plan_id=plan_id,
|
|
plan_status=updated_plan.status.value if hasattr(updated_plan.status, "value") else updated_plan.status,
|
|
generation_task_id=gen_task.id,
|
|
clip_count=clip_count,
|
|
)
|
|
except HTTPException:
|
|
raise
|
|
except Exception as _e:
|
|
logger.exception("触发剪辑计划生成失败: plan_id=%s", plan_id)
|
|
try:
|
|
svc.transition_status(plan_id, EditPlanStatus.FAILED)
|
|
except Exception:
|
|
logger.warning("标记计划失败状态时异常: plan_id=%s", plan_id)
|
|
raise HTTPException(
|
|
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
|
detail="生成失败,请稍后重试",
|
|
) from _e
|
|
|
|
|
|
@router.get(
|
|
"/{plan_id}/generation-status",
|
|
response_model=EditPlanGenerationStatusResponse,
|
|
)
|
|
def get_generation_status(
|
|
plan_id: str,
|
|
db: Session = Depends(get_db_session),
|
|
current_user: AuthenticatedUser = Depends(get_current_user),
|
|
project_repository: Any = Depends(get_project_repository),
|
|
) -> EditPlanGenerationStatusResponse:
|
|
"""查询剪辑计划生成进度"""
|
|
svc = EditPlanService(db)
|
|
try:
|
|
gen_status = svc.get_generation_status(plan_id)
|
|
except ValueError as exc:
|
|
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=str(exc)) from exc
|
|
|
|
plan = gen_status["plan"]
|
|
if plan.project_id:
|
|
check_project_access(plan.project_id, current_user.user.id, project_repository)
|
|
clips = gen_status["clips"]
|
|
|
|
clip_items = [
|
|
ClipStatusItem(
|
|
clip_id=c.id,
|
|
clip_type=c.clip_type,
|
|
order=c.order,
|
|
status=c.status.value if hasattr(c.status, "value") else c.status,
|
|
asset_id=c.asset_id or "",
|
|
text_content=c.text_content or "",
|
|
duration=c.duration,
|
|
)
|
|
for c in clips
|
|
]
|
|
|
|
# 从 plan.config 中取渲染结果 URL
|
|
video_url = (plan.config or {}).get("rendered_url", "")
|
|
# 从 gen_status 中取进度、错误信息、任务状态
|
|
progress = gen_status.get("progress", 0.0)
|
|
error_message = gen_status.get("error_message", "")
|
|
gen_task_status = gen_status.get("generation_task_status")
|
|
# 如果计划已完成但进度还是0,补100
|
|
plan_status_val = plan.status.value if hasattr(plan.status, "value") else plan.status
|
|
if plan_status_val == "completed" and progress < 100:
|
|
progress = 100.0
|
|
|
|
return EditPlanGenerationStatusResponse(
|
|
plan_id=plan_id,
|
|
plan_status=plan_status_val,
|
|
generation_task_id=gen_status["generation_task_id"],
|
|
generation_task_status=gen_task_status,
|
|
progress=progress,
|
|
video_url=video_url,
|
|
error_message=error_message,
|
|
clips=clip_items,
|
|
)
|
|
|
|
|
|
@router.get(
|
|
"/{plan_id}/generations",
|
|
response_model=EditPlanGenerationsResponse,
|
|
)
|
|
def list_plan_generations(
|
|
plan_id: str,
|
|
db: Session = Depends(get_db_session),
|
|
current_user: AuthenticatedUser = Depends(get_current_user),
|
|
project_repository: Any = Depends(get_project_repository),
|
|
) -> EditPlanGenerationsResponse:
|
|
"""查询剪辑计划关联的所有生成记录"""
|
|
svc = EditPlanService(db)
|
|
plan = svc.get_plan_or_raise(plan_id)
|
|
if plan.project_id:
|
|
check_project_access(plan.project_id, current_user.user.id, project_repository)
|
|
|
|
from app.schemas.generation_task import GenerationTaskResponse
|
|
|
|
gen_task_repo = SQLAlchemyGenerationTaskRepository(db)
|
|
tasks = gen_task_repo.list_by_source_edit_plan(plan_id)
|
|
items = [
|
|
GenerationTaskResponse(
|
|
id=t.id,
|
|
project_id=t.project_id,
|
|
asset_library_id=t.asset_library_id,
|
|
strategy_id=t.strategy_id,
|
|
voice_library_id=t.voice_library_id,
|
|
template_id=t.template_id,
|
|
asset_ids=t.asset_ids,
|
|
title_ids=t.title_ids,
|
|
voice_ids=t.voice_ids,
|
|
source_edit_plan_id=t.source_edit_plan_id or "",
|
|
status=t.status.value if hasattr(t.status, "value") else t.status,
|
|
progress=t.progress,
|
|
result_count=t.result_count,
|
|
error_message=t.error_message,
|
|
)
|
|
for t in tasks
|
|
]
|
|
return EditPlanGenerationsResponse(items=items, total=len(items))
|