9c99c9ea96
CI/CD Pipeline / Validate Code Quality And Tests (push) Failing after 8s
CI/CD Pipeline / Validate Code Quality And Tests (pull_request) Failing after 8s
CI/CD Pipeline / Production Browser E2E (pull_request) Failing after 1648h53m1s
CI/CD Pipeline / Staging API Integration Tests (pull_request) Failing after 1648h53m3s
CI/CD Pipeline / Staging E2E Tests (pull_request) Failing after 1648h53m3s
CI/CD Pipeline / Build & Push Staging (Watchtower auto-deploy) (pull_request) Failing after 1648h53m5s
CI/CD Pipeline / Production Browser E2E (push) Failing after 1648h53m29s
CI/CD Pipeline / Build Production Runtime Images (push) Failing after 1648h53m34s
CI/CD Pipeline / Deploy Production (push) Failing after 1648h53m32s
CI/CD Pipeline / Build & Push Staging (Watchtower auto-deploy) (push) Failing after 1648h53m34s
CI/CD Pipeline / Staging API Integration Tests (push) Failing after 1648h53m32s
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Deploy Production (pull_request) Failing after 1649h24m33s
CI/CD Pipeline / Build Production Runtime Images (pull_request) Failing after 1649h24m35s
CI/CD Pipeline / Frontend Lint (push) Has been skipped
CI/CD Pipeline / Staging E2E Tests (push) Failing after 1649h25m2s
核心变更: 1. 提取共享工具模块(ffmpeg_utils / oss_helpers / dedup_helpers) 2. 实现 UnifiedRenderService — 按 clip_type/config.role 分组为图层再合成 3. 重构 render_edit_plan() 使用 UnifiedRenderService(替换 concat demuxer) 4. 重构 generate_video() 使用 UnifiedRenderService(替换 EditingModeProcessor) 5. 集成 VideoDeduplicator 查重 6. 补 23 个单元测试 + 14 个四模式集成测试 + 6 个全链路测试 图层分组算法: main → main (z=0) main+config.role=b_roll → broll (z=0) overlay → overlay (z=1) background → background (z=-1) corner_voice → corner_voice (z=1) b_roll → broll (z=0) intro/outro → main (z=0) 合成流程:每个 clip 预处理 → 同层 xfade 串联 → overlay 合成 → 音频混入
490 lines
17 KiB
Python
Executable File
490 lines
17 KiB
Python
Executable File
"""
|
|
视频生成任务 — 使用 UnifiedRenderService 统一渲染引擎.
|
|
|
|
支持四种剪辑模式:一镜到底、画中画、口播、口播+画中画。
|
|
模式差异体现在虚拟剪辑计划的 clip_type 分布上,渲染引擎不判断模式。
|
|
|
|
模式 → clip_type 映射:
|
|
ONE_TAKE: N 个 main clips
|
|
PIP: 1 main + N-1 overlay
|
|
VOICE_OVER: N 个 main(config.role=b_roll)
|
|
VOICE_PIP: 1 background + 1 corner_voice + N-2 b_roll
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import os
|
|
import tempfile
|
|
from dataclasses import dataclass, field
|
|
from pathlib import Path
|
|
from typing import Any, Optional
|
|
|
|
from worker_app.celery_app import celery_app
|
|
from worker_app.db import SessionLocal
|
|
|
|
OUTPUT_WIDTH = 1280
|
|
OUTPUT_HEIGHT = 720
|
|
OUTPUT_FPS = 25.0
|
|
OUTPUT_DURATION_SECONDS = 5.0
|
|
GENERATED_FILES_DIR = Path(os.getenv("GENERATED_FILES_DIR", "/app/generated"))
|
|
GENERATED_FILES_URL_PREFIX = os.getenv("GENERATED_FILES_URL_PREFIX", "/generated-files")
|
|
PUBLIC_API_BASE_URL = os.getenv("PUBLIC_API_BASE_URL", "https://api.xiaoxiajianji.com").rstrip("/")
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
# ── 状态更新辅助函数 ──────────────────────────────────────────────────────────
|
|
|
|
|
|
def _update_task_status(task_id: str, status_action: str, **kwargs) -> bool:
|
|
"""更新 GenerationTask 状态(独立 session,异常不向外抛出)。
|
|
|
|
Args:
|
|
task_id: 任务 ID
|
|
status_action: 状态动作名,如 "mark_processing" / "mark_completed" / "mark_failed"
|
|
**kwargs: 传递给对应方法的参数
|
|
|
|
Returns:
|
|
True 表示更新成功,False 表示更新失败
|
|
"""
|
|
try:
|
|
from packages.adapters.sqlalchemy_impl.generation_task_repository import (
|
|
SQLAlchemyGenerationTaskRepository,
|
|
)
|
|
|
|
session = SessionLocal()
|
|
try:
|
|
repo = SQLAlchemyGenerationTaskRepository(session)
|
|
task = repo.get(task_id)
|
|
if task is None:
|
|
logger.warning("更新任务状态失败:任务不存在 task_id=%s", task_id)
|
|
return False
|
|
|
|
action = getattr(task, status_action, None)
|
|
if action is None:
|
|
logger.warning("未知的状态动作: %s", status_action)
|
|
return False
|
|
|
|
action(**kwargs)
|
|
repo.update(task)
|
|
logger.info("GenerationTask 状态更新成功: task_id=%s action=%s", task_id, status_action)
|
|
return True
|
|
finally:
|
|
session.close()
|
|
except Exception as e:
|
|
logger.error(
|
|
"更新 GenerationTask 状态异常: task_id=%s action=%s error=%s",
|
|
task_id,
|
|
status_action,
|
|
e,
|
|
exc_info=True,
|
|
)
|
|
return False
|
|
|
|
|
|
# ── 共享工具模块导入 ──────────────────────────────────────────────────────────
|
|
|
|
from video_processing.ffmpeg_utils import FFMPEG_BIN, run_ffmpeg, probe_duration
|
|
from video_processing.oss_helpers import (
|
|
download_asset,
|
|
oss_bucket,
|
|
upload_to_oss,
|
|
)
|
|
from video_processing.dedup_helpers import create_video_record_and_dedup
|
|
from video_processing.unified_render_service import UnifiedRenderService
|
|
|
|
|
|
# ── 虚拟 Plan / Clip(内存中构建,不写数据库) ────────────────────────────────
|
|
|
|
|
|
@dataclass
|
|
class _VirtualPlan:
|
|
"""内存中的虚拟剪辑计划,供 UnifiedRenderService 使用。"""
|
|
|
|
id: str
|
|
name: str = ""
|
|
|
|
|
|
@dataclass
|
|
class _VirtualClip:
|
|
"""内存中的虚拟剪辑片段,供 UnifiedRenderService 使用。"""
|
|
|
|
id: str
|
|
plan_id: str = ""
|
|
clip_type: str = "main"
|
|
order: int = 0
|
|
asset_id: str = ""
|
|
text_content: str = ""
|
|
start_time: float = 0.0
|
|
duration: float = 0.0
|
|
transition_effect: str = "cut"
|
|
status: str = "ready"
|
|
config: dict[str, Any] = field(default_factory=dict)
|
|
|
|
|
|
def _build_plan_and_clips_from_task(
|
|
task_id: str,
|
|
downloaded_paths: list[Path],
|
|
mode: str,
|
|
) -> tuple[_VirtualPlan, list[_VirtualClip], dict[str, Path]]:
|
|
"""根据模式和下载的素材路径,构建虚拟 plan + clips + asset_path_map。
|
|
|
|
模式 → clip_type 映射:
|
|
ONE_TAKE: N 个 main clips
|
|
PIP: 1 main + N-1 overlay
|
|
VOICE_OVER: N 个 main(config.role=b_roll)
|
|
VOICE_PIP: 1 background + 1 corner_voice + N-2 b_roll
|
|
|
|
Returns:
|
|
(virtual_plan, virtual_clips, asset_path_map)
|
|
"""
|
|
plan = _VirtualPlan(id=task_id, name=f"Generated-{task_id[:8]}")
|
|
|
|
# 为每个下载路径生成合成 asset_id
|
|
asset_path_map: dict[str, Path] = {}
|
|
path_to_asset_id: dict[Path, str] = {}
|
|
for i, p in enumerate(downloaded_paths):
|
|
asset_id = f"gen_{task_id[:8]}_{i:03d}{p.suffix or '.mp4'}"
|
|
asset_path_map[asset_id] = p
|
|
path_to_asset_id[p] = asset_id
|
|
|
|
clips: list[_VirtualClip] = []
|
|
n = len(downloaded_paths)
|
|
|
|
if mode == "pip":
|
|
# 1 main + N-1 overlay
|
|
for i, p in enumerate(downloaded_paths):
|
|
clip_type = "main" if i == 0 else "overlay"
|
|
clips.append(_VirtualClip(
|
|
id=f"vc_{i:03d}",
|
|
plan_id=task_id,
|
|
clip_type=clip_type,
|
|
order=i,
|
|
asset_id=path_to_asset_id[p],
|
|
))
|
|
elif mode == "voice_over":
|
|
# N 个 main(config.role=b_roll)
|
|
for i, p in enumerate(downloaded_paths):
|
|
clips.append(_VirtualClip(
|
|
id=f"vc_{i:03d}",
|
|
plan_id=task_id,
|
|
clip_type="main",
|
|
order=i,
|
|
asset_id=path_to_asset_id[p],
|
|
config={"role": "b_roll"},
|
|
))
|
|
elif mode == "voice_pip":
|
|
# 1 background + 1 corner_voice + N-2 b_roll
|
|
for i, p in enumerate(downloaded_paths):
|
|
if i == 0:
|
|
clip_type = "background"
|
|
elif i == 1:
|
|
clip_type = "corner_voice"
|
|
else:
|
|
clip_type = "b_roll"
|
|
clips.append(_VirtualClip(
|
|
id=f"vc_{i:03d}",
|
|
plan_id=task_id,
|
|
clip_type=clip_type,
|
|
order=i,
|
|
asset_id=path_to_asset_id[p],
|
|
))
|
|
else:
|
|
# ONE_TAKE (default): N 个 main clips
|
|
for i, p in enumerate(downloaded_paths):
|
|
clips.append(_VirtualClip(
|
|
id=f"vc_{i:03d}",
|
|
plan_id=task_id,
|
|
clip_type="main",
|
|
order=i,
|
|
asset_id=path_to_asset_id[p],
|
|
))
|
|
|
|
return plan, clips, asset_path_map
|
|
|
|
|
|
def _create_fallback_clip(output_path: Path, title: str) -> None:
|
|
"""创建 fallback 视频(无素材时)"""
|
|
safe_title = title.replace(":", "\\:").replace("'", "\\'")[:80]
|
|
run_ffmpeg(
|
|
[
|
|
FFMPEG_BIN,
|
|
"-y",
|
|
"-f",
|
|
"lavfi",
|
|
"-i",
|
|
f"color=c=#111827:s={OUTPUT_WIDTH}x{OUTPUT_HEIGHT}:d={OUTPUT_DURATION_SECONDS}:r={int(OUTPUT_FPS)}",
|
|
"-vf",
|
|
f"drawtext=text='{safe_title}':fontcolor=white:fontsize=48:x=(w-text_w)/2:y=(h-text_h)/2",
|
|
"-c:v",
|
|
"libx264",
|
|
"-pix_fmt",
|
|
"yuv420p",
|
|
"-movflags",
|
|
"+faststart",
|
|
str(output_path),
|
|
]
|
|
)
|
|
|
|
|
|
def _mux_audio_track(video_path: Path, audio_path: str, output_path: Path) -> None:
|
|
"""将音频轨混入已渲染的视频(后处理步骤)。
|
|
|
|
使用 FFmpeg 将视频和音频合并,视频时长为准,音频不足则循环,
|
|
音频过长则截断。
|
|
"""
|
|
command = [
|
|
FFMPEG_BIN,
|
|
"-y",
|
|
"-i", str(video_path),
|
|
"-i", audio_path,
|
|
"-c:v", "copy",
|
|
"-c:a", "aac",
|
|
"-b:a", "192k",
|
|
"-shortest",
|
|
"-map", "0:v:0",
|
|
"-map", "1:a:0",
|
|
"-movflags", "+faststart",
|
|
str(output_path),
|
|
]
|
|
run_ffmpeg(command)
|
|
|
|
|
|
def _download_voice_asset(voice_library_id: str, local_path: Path) -> bool:
|
|
"""下载配音文件"""
|
|
if not voice_library_id:
|
|
return False
|
|
storage_key = f"voice/{voice_library_id}.mp3"
|
|
return download_asset(storage_key, local_path)
|
|
|
|
|
|
def _download_library_assets(
|
|
asset_library_id: str,
|
|
temp_path: Path,
|
|
video_extensions: tuple = (".mp4", ".mov", ".avi", ".mkv", ".webm"),
|
|
asset_ids: list[str] | None = None,
|
|
) -> list[Path]:
|
|
"""从素材库下载视频素材。
|
|
|
|
Args:
|
|
asset_library_id: 素材库 ID
|
|
temp_path: 临时目录路径
|
|
video_extensions: 支持的视频扩展名
|
|
asset_ids: 指定素材 ID 列表,为空则下载全部 ready 视频素材
|
|
|
|
Returns:
|
|
下载成功的视频文件 Path 列表
|
|
"""
|
|
try:
|
|
from packages.adapters.sqlalchemy_impl.models import AssetModel
|
|
|
|
session = SessionLocal()
|
|
try:
|
|
query = session.query(AssetModel).filter(
|
|
AssetModel.asset_library_id == asset_library_id,
|
|
AssetModel.status == "ready",
|
|
AssetModel.file_type.in_(["video", "video/mp4", "video/quicktime"]),
|
|
)
|
|
if asset_ids:
|
|
query = query.filter(AssetModel.id.in_(asset_ids))
|
|
assets = query.order_by(AssetModel.created_at).all()
|
|
|
|
if not assets:
|
|
logger.info("No video assets found in library %s", asset_library_id)
|
|
return []
|
|
|
|
downloaded: list[Path] = []
|
|
for i, asset in enumerate(assets):
|
|
storage_key = asset.file_url if asset.file_url else None
|
|
if not storage_key:
|
|
continue
|
|
|
|
ext = Path(storage_key).suffix or ".mp4"
|
|
local_file = temp_path / f"asset_{i:03d}_{asset.id}{ext}"
|
|
if download_asset(storage_key, local_file):
|
|
downloaded.append(local_file)
|
|
logger.info("Downloaded asset: %s -> %s", asset.name, local_file)
|
|
else:
|
|
logger.warning("Failed to download asset: %s", asset.name)
|
|
|
|
return downloaded
|
|
finally:
|
|
session.close()
|
|
except Exception as e:
|
|
logger.error("Error downloading library assets: %s", e)
|
|
return []
|
|
|
|
|
|
# ── Celery Task ──────────────────────────────────────────────────────────────
|
|
|
|
|
|
@celery_app.task(bind=True, name="worker.generate_video", max_retries=2)
|
|
def generate_video(self, task_id: str) -> dict:
|
|
"""生成视频任务 — 使用 UnifiedRenderService 统一渲染。
|
|
|
|
流程:
|
|
1. 加载 GenerationTask 信息
|
|
2. 从素材库下载视频素材
|
|
3. 根据模式构建虚拟 plan + clips
|
|
4. 使用 UnifiedRenderService 渲染
|
|
5. 如有配音,后处理混音
|
|
6. 上传 OSS + 查重
|
|
7. 更新 GenerationTask 状态
|
|
|
|
Args:
|
|
task_id: 任务 ID(从数据库加载完整任务信息)
|
|
|
|
Returns:
|
|
生成结果字典
|
|
"""
|
|
from packages.domain import EditingMode
|
|
|
|
logger.info("开始生成视频任务: task_id=%s", task_id)
|
|
|
|
# 从数据库加载任务信息
|
|
session = SessionLocal()
|
|
try:
|
|
from packages.adapters.sqlalchemy_impl.generation_task_repository import (
|
|
SQLAlchemyGenerationTaskRepository,
|
|
)
|
|
|
|
task_repo = SQLAlchemyGenerationTaskRepository(session)
|
|
gen_task = task_repo.get(task_id)
|
|
if gen_task is None:
|
|
logger.error("生成任务不存在: task_id=%s", task_id)
|
|
return {"status": "failed", "error": f"generation task {task_id} not found"}
|
|
project_id = gen_task.project_id
|
|
asset_library_id = gen_task.asset_library_id
|
|
voice_library_id = gen_task.voice_library_id or ""
|
|
mode = gen_task.strategy_id or "one_take"
|
|
task_asset_ids = list(gen_task.asset_ids or [])
|
|
batch_id = getattr(gen_task, "batch_id", "") or ""
|
|
finally:
|
|
session.close()
|
|
|
|
# 标记任务为 running
|
|
_update_task_status(task_id, "mark_processing")
|
|
|
|
try:
|
|
editing_mode = EditingMode(mode)
|
|
except ValueError:
|
|
editing_mode = EditingMode.ONE_TAKE
|
|
|
|
output_name = f"generated-{task_id}.mp4"
|
|
storage_key = f"generated/projects/{project_id}/tasks/{task_id}/{output_name}"
|
|
|
|
try:
|
|
with tempfile.TemporaryDirectory(prefix="xiaoxia-generation-") as temp_dir:
|
|
temp_path = Path(temp_dir)
|
|
output_path = temp_path / output_name
|
|
|
|
# 1. 从素材库下载视频素材
|
|
downloaded_videos = _download_library_assets(
|
|
asset_library_id, temp_path, asset_ids=task_asset_ids or None
|
|
)
|
|
|
|
# 2. 下载配音(如有)
|
|
audio_path: str | None = None
|
|
if voice_library_id:
|
|
local_audio = temp_path / "voice.mp3"
|
|
if _download_voice_asset(voice_library_id, local_audio):
|
|
audio_path = str(local_audio)
|
|
|
|
# 3. 渲染
|
|
if downloaded_videos:
|
|
# 构建虚拟 plan + clips + asset_path_map
|
|
virtual_plan, virtual_clips, asset_path_map = _build_plan_and_clips_from_task(
|
|
task_id=task_id,
|
|
downloaded_paths=downloaded_videos,
|
|
mode=editing_mode.value,
|
|
)
|
|
|
|
# 使用 UnifiedRenderService 渲染
|
|
render_service = UnifiedRenderService(
|
|
plan=virtual_plan,
|
|
clips=virtual_clips,
|
|
asset_path_map=asset_path_map,
|
|
work_dir=temp_path,
|
|
output_width=OUTPUT_WIDTH,
|
|
output_height=OUTPUT_HEIGHT,
|
|
output_fps=int(OUTPUT_FPS),
|
|
)
|
|
render_result = render_service.render()
|
|
|
|
# 4. 如有配音,后处理混音
|
|
if audio_path:
|
|
final_path = temp_path / f"final-{task_id}.mp4"
|
|
try:
|
|
_mux_audio_track(render_result.output_path, audio_path, final_path)
|
|
# 混音成功,使用混音后的文件
|
|
output_path = final_path
|
|
except Exception as mux_err:
|
|
logger.warning("音频混合失败,使用无音频版本: %s", mux_err)
|
|
output_path = render_result.output_path
|
|
else:
|
|
output_path = render_result.output_path
|
|
else:
|
|
# 无素材,生成 fallback 视频
|
|
_create_fallback_clip(output_path, f"Generated Video {task_id[:8]}")
|
|
|
|
file_size = output_path.stat().st_size
|
|
duration = probe_duration(output_path)
|
|
|
|
# 5. 上传到 OSS
|
|
bucket = oss_bucket()
|
|
if bucket:
|
|
try:
|
|
bucket.put_object_from_file(storage_key, str(output_path))
|
|
except Exception as oss_err:
|
|
logger.warning("OSS upload failed: %s", oss_err)
|
|
|
|
# 构建视频 URL
|
|
if bucket:
|
|
file_url = f"{PUBLIC_API_BASE_URL}/{storage_key}"
|
|
else:
|
|
file_url = f"{GENERATED_FILES_URL_PREFIX}/{task_id}/{output_name}"
|
|
|
|
# 6. 创建 GeneratedVideo 记录 + 查重
|
|
dedup_session = SessionLocal()
|
|
try:
|
|
video_count = create_video_record_and_dedup(
|
|
generation_task_id=task_id,
|
|
project_id=project_id,
|
|
batch_id=batch_id,
|
|
file_url=file_url,
|
|
file_size=file_size,
|
|
duration=duration,
|
|
video_path=str(output_path),
|
|
mode=editing_mode.value,
|
|
session=dedup_session,
|
|
)
|
|
finally:
|
|
dedup_session.close()
|
|
|
|
# 7. 标记任务为 completed
|
|
_update_task_status(task_id, "mark_completed", result_count=video_count or 1)
|
|
|
|
logger.info("视频生成完成: task_id=%s duration=%.2fs file_size=%d", task_id, duration, file_size)
|
|
|
|
return {
|
|
"status": "completed",
|
|
"task_id": task_id,
|
|
"output_path": str(output_path),
|
|
"file_size": file_size,
|
|
"duration": duration,
|
|
"width": OUTPUT_WIDTH,
|
|
"height": OUTPUT_HEIGHT,
|
|
"mode": editing_mode.value,
|
|
}
|
|
except Exception as error:
|
|
logger.error("Video generation failed: %s", error, exc_info=True)
|
|
_update_task_status(task_id, "mark_failed", error_message=str(error))
|
|
return {
|
|
"status": "failed",
|
|
"task_id": task_id,
|
|
"error": str(error),
|
|
}
|
|
|
|
|