fdcf48103e
CI/CD Pipeline / Validate Code Quality And Tests (push) Successful in 1m27s
CI/CD Pipeline / Frontend Lint (push) Successful in 2m4s
CI/CD Pipeline / Build Production Runtime Images (push) Has been skipped
CI/CD Pipeline / Deploy Production (push) Has been skipped
CI/CD Pipeline / Production Browser E2E (push) Has been skipped
CI/CD Pipeline / Integration Tests (push) Successful in 1m44s
CI/CD Pipeline / Build & Push Staging (Watchtower auto-deploy) (push) Successful in 9m39s
CI/CD Pipeline / Staging E2E Tests (push) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (push) Has been skipped
Phase 3 unified render with audio mixing + gray scale observation metrics
300 lines
10 KiB
Python
Executable File
300 lines
10 KiB
Python
Executable File
"""统一渲染引擎适配层 — Phase 2.
|
||
|
||
将 EditPlan + EditPlanClips(来自 DB)适配为 UnifiedRenderService 的输入格式,
|
||
封装素材下载、渲染执行、结果上传的完整流程。
|
||
|
||
职责:
|
||
1. 从 DB 读取 EditPlan + EditPlanClips
|
||
2. 下载素材到本地,构建 asset_path_map
|
||
3. 调用 UnifiedRenderService 执行渲染
|
||
4. 上传渲染结果到 OSS
|
||
5. 支持进度回调(对接 JobService)
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
import tempfile
|
||
from dataclasses import dataclass
|
||
from pathlib import Path
|
||
from typing import Any, Callable
|
||
|
||
from sqlalchemy.orm import Session
|
||
from video_processing.oss_helpers import download_asset, upload_to_oss
|
||
from video_processing.unified_render_service import RenderResult, UnifiedRenderService
|
||
|
||
from packages.adapters.sqlalchemy_impl.edit_plan_clip_repository import SQLAlchemyEditPlanClipRepository
|
||
from packages.adapters.sqlalchemy_impl.edit_plan_repository import SQLAlchemyEditPlanRepository
|
||
from packages.domain.edit_plan import EditPlan, EditPlanStatus
|
||
from packages.domain.edit_plan_clip import EditPlanClip, EditPlanClipStatus
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
|
||
# ── 数据结构 ──────────────────────────────────────────────────────────────────
|
||
|
||
|
||
@dataclass
|
||
class RenderAdapterResult:
|
||
"""渲染适配结果。"""
|
||
|
||
success: bool
|
||
output_url: str = ""
|
||
output_path: Path | None = None
|
||
duration: float = 0.0
|
||
file_size: int = 0
|
||
width: int = 0
|
||
height: int = 0
|
||
clip_count: int = 0
|
||
error_message: str = ""
|
||
|
||
|
||
ProgressCallback = Callable[[float, str], None]
|
||
"""进度回调:(progress_0_100, stage_description) → None"""
|
||
|
||
|
||
# ── 适配层主体 ────────────────────────────────────────────────────────────────
|
||
|
||
|
||
class RenderAdapter:
|
||
"""统一渲染引擎适配层。
|
||
|
||
桥接 EditPlan 领域模型与 UnifiedRenderService 图层模型。
|
||
|
||
用法::
|
||
|
||
adapter = RenderAdapter(db)
|
||
result = adapter.render_plan(
|
||
plan_id=plan_id,
|
||
job_id=job_id,
|
||
progress_cb=lambda p, s: job_service.update_progress(job_id, p, s),
|
||
)
|
||
"""
|
||
|
||
def __init__(self, db: Session) -> None:
|
||
self._db = db
|
||
self._plan_repo = SQLAlchemyEditPlanRepository(db)
|
||
self._clip_repo = SQLAlchemyEditPlanClipRepository(db)
|
||
|
||
# ── 公开方法 ──────────────────────────────────────────────────────────
|
||
|
||
def render_plan(
|
||
self,
|
||
plan_id: str,
|
||
*,
|
||
job_id: str = "",
|
||
work_dir: Path | None = None,
|
||
progress_cb: ProgressCallback | None = None,
|
||
) -> RenderAdapterResult:
|
||
"""渲染一个 EditPlan。
|
||
|
||
完整流程:
|
||
1. 加载计划与片段
|
||
2. 下载素材
|
||
3. 执行统一渲染
|
||
4. 上传结果
|
||
|
||
Args:
|
||
plan_id: EditPlan ID
|
||
job_id: 关联的 Job ID(用于结果存储路径)
|
||
work_dir: 工作目录,不传则使用临时目录
|
||
progress_cb: 进度回调函数
|
||
|
||
Returns:
|
||
RenderAdapterResult
|
||
"""
|
||
temp_dir = None
|
||
try:
|
||
# 0. 准备工作目录
|
||
if work_dir is None:
|
||
temp_dir = tempfile.mkdtemp(prefix="render_")
|
||
work_dir = Path(temp_dir)
|
||
work_dir.mkdir(parents=True, exist_ok=True)
|
||
|
||
self._report_progress(progress_cb, 5.0, "加载剪辑计划")
|
||
|
||
# 1. 加载计划与片段
|
||
plan = self._plan_repo.get(plan_id)
|
||
if plan is None:
|
||
return RenderAdapterResult(
|
||
success=False,
|
||
error_message=f"剪辑计划不存在: {plan_id}",
|
||
)
|
||
|
||
clips = self._clip_repo.list_by_plan(plan_id, skip=0, limit=10000)
|
||
ready_clips = [c for c in clips if c.status == EditPlanClipStatus.READY and c.asset_id]
|
||
ready_clips.sort(key=lambda c: c.order)
|
||
|
||
if not ready_clips:
|
||
return RenderAdapterResult(
|
||
success=False,
|
||
error_message="没有可渲染的就绪片段",
|
||
clip_count=0,
|
||
)
|
||
|
||
logger.info(
|
||
"开始渲染: plan_id=%s job_id=%s ready_clips=%d engine=unified",
|
||
plan_id,
|
||
job_id,
|
||
len(ready_clips),
|
||
)
|
||
|
||
self._report_progress(progress_cb, 15.0, f"下载素材({len(ready_clips)} 个)")
|
||
|
||
# 2. 下载素材
|
||
asset_path_map = self._download_assets(ready_clips, work_dir)
|
||
if not asset_path_map:
|
||
return RenderAdapterResult(
|
||
success=False,
|
||
error_message="所有素材下载失败",
|
||
clip_count=len(ready_clips),
|
||
)
|
||
|
||
self._report_progress(progress_cb, 40.0, "执行视频渲染")
|
||
|
||
# 3. 执行统一渲染
|
||
render_svc = UnifiedRenderService(
|
||
plan=plan,
|
||
clips=ready_clips,
|
||
asset_path_map=asset_path_map,
|
||
work_dir=work_dir,
|
||
)
|
||
result = render_svc.render()
|
||
|
||
self._report_progress(progress_cb, 80.0, "上传渲染结果")
|
||
|
||
# 4. 上传结果
|
||
storage_key = f"rendered/{plan_id}/{job_id or plan_id}.mp4"
|
||
output_url = upload_to_oss(result.output_path, storage_key)
|
||
|
||
self._report_progress(progress_cb, 100.0, "渲染完成")
|
||
|
||
logger.info(
|
||
"[render-adapter] render success: plan_id=%s job_id=%s engine=unified "
|
||
"duration=%.2fs file_size=%d resolution=%dx%d clip_count=%d",
|
||
plan_id,
|
||
job_id,
|
||
result.duration,
|
||
result.file_size,
|
||
result.width,
|
||
result.height,
|
||
len(ready_clips),
|
||
)
|
||
|
||
return RenderAdapterResult(
|
||
success=True,
|
||
output_url=output_url or "",
|
||
output_path=result.output_path,
|
||
duration=result.duration,
|
||
file_size=result.file_size,
|
||
width=result.width,
|
||
height=result.height,
|
||
clip_count=len(ready_clips),
|
||
)
|
||
|
||
except Exception as exc:
|
||
logger.exception(
|
||
"[render-adapter] render failed: plan_id=%s job_id=%s engine=unified error=%s",
|
||
plan_id,
|
||
job_id,
|
||
str(exc)[:200],
|
||
)
|
||
return RenderAdapterResult(
|
||
success=False,
|
||
error_message=str(exc)[:500],
|
||
)
|
||
finally:
|
||
# 清理临时目录
|
||
if temp_dir:
|
||
import shutil
|
||
|
||
try:
|
||
shutil.rmtree(temp_dir, ignore_errors=True)
|
||
except Exception:
|
||
pass
|
||
|
||
def validate_plan(self, plan_id: str) -> tuple[bool, list[str], list[str], int, int]:
|
||
"""校验计划是否可渲染(兼容 VideoComposeService.validate_compose 接口)。
|
||
|
||
Returns:
|
||
(valid, errors, warnings, ready_clip_count, total_clip_count)
|
||
"""
|
||
errors: list[str] = []
|
||
warnings: list[str] = []
|
||
|
||
plan = self._plan_repo.get(plan_id)
|
||
if plan is None:
|
||
return False, [f"剪辑计划不存在: {plan_id}"], [], 0, 0
|
||
|
||
if plan.status not in (EditPlanStatus.EDITING, EditPlanStatus.RENDERING):
|
||
errors.append(f"计划状态不正确,需要 editing 或 rendering,当前: {plan.status}")
|
||
|
||
clips = self._clip_repo.list_by_plan(plan_id, skip=0, limit=10000)
|
||
if not clips:
|
||
errors.append("计划没有任何片段")
|
||
return False, errors, warnings, 0, 0
|
||
|
||
clips.sort(key=lambda c: c.order)
|
||
|
||
ready_count = 0
|
||
pending_count = 0
|
||
no_asset_count = 0
|
||
|
||
for clip in clips:
|
||
if clip.status == EditPlanClipStatus.READY:
|
||
ready_count += 1
|
||
if not clip.asset_id:
|
||
errors.append(f"片段 {clip.id} (order={clip.order}) 没有分配素材")
|
||
no_asset_count += 1
|
||
elif clip.status == EditPlanClipStatus.PENDING:
|
||
pending_count += 1
|
||
elif clip.status == EditPlanClipStatus.FAILED:
|
||
warnings.append(f"片段 {clip.id} (order={clip.order}) 状态为 failed,已跳过")
|
||
|
||
if ready_count == 0:
|
||
errors.append("没有就绪(ready)的片段可以合成")
|
||
|
||
if pending_count > 0:
|
||
warnings.append(f"有 {pending_count} 个片段仍处于 pending 状态")
|
||
|
||
return len(errors) == 0, errors, warnings, ready_count, len(clips)
|
||
|
||
# ── 内部方法 ──────────────────────────────────────────────────────────
|
||
|
||
@staticmethod
|
||
def _report_progress(progress_cb: ProgressCallback | None, progress: float, stage: str) -> None:
|
||
"""上报进度。"""
|
||
if progress_cb is not None:
|
||
try:
|
||
progress_cb(progress, stage)
|
||
except Exception:
|
||
logger.exception("进度回调失败")
|
||
|
||
@staticmethod
|
||
def _download_assets(clips: list[EditPlanClip], work_dir: Path) -> dict[str, Path]:
|
||
"""下载片段素材到本地,返回 asset_id → local_path 映射。
|
||
|
||
只保留下载成功的素材。
|
||
"""
|
||
asset_dir = work_dir / "assets"
|
||
asset_dir.mkdir(exist_ok=True)
|
||
|
||
asset_path_map: dict[str, Path] = {}
|
||
|
||
for clip in clips:
|
||
asset_id = clip.asset_id
|
||
if not asset_id:
|
||
continue
|
||
|
||
# 生成安全的本地文件名
|
||
safe_name = f"clip_{clip.order:04d}_{abs(hash(asset_id)) % 100000:05d}.mp4"
|
||
local_path = asset_dir / safe_name
|
||
|
||
if download_asset(asset_id, local_path):
|
||
asset_path_map[asset_id] = local_path
|
||
logger.debug("素材下载成功: clip_id=%s asset_id=%s", clip.id, asset_id[:60])
|
||
else:
|
||
logger.warning("素材下载失败: clip_id=%s asset_id=%s", clip.id, asset_id[:60])
|
||
|
||
return asset_path_map
|