Compare commits

...

9 Commits

Author SHA1 Message Date
Deploy Bot cb26ccb147 merge: develop(FFmpeg超时保护) into stream-copy分支
CI/CD Pipeline / Validate Code Quality And Tests (pull_request) Failing after 1m57s
CI/CD Pipeline / Frontend Lint (pull_request) Successful in 1m56s
CI/CD Pipeline / Build & Push Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Build Production Runtime Images (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Integration Tests (pull_request) Successful in 1m57s
2026-07-13 11:09:09 +08:00
CI Bot b094346bf6 fix(worker): 修复 flake8 错误 - F401/F841/E501
CI/CD Pipeline / Validate Code Quality And Tests (pull_request) Failing after 1m40s
CI/CD Pipeline / Frontend Lint (pull_request) Successful in 1m54s
CI/CD Pipeline / Build & Push Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Build Production Runtime Images (pull_request) Has been skipped
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Integration Tests (pull_request) Successful in 2m0s
2026-07-13 11:08:21 +08:00
xiaoxia a0cac1b75d fix(worker): FFmpeg超时保护 - 防止渲染hang住导致worker永久阻塞 (#242)
CI/CD Pipeline / Validate Code Quality And Tests (push) Successful in 2m1s
CI/CD Pipeline / Frontend Lint (push) Successful in 2m10s
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 2m3s
CI/CD Pipeline / Build & Push Staging (Watchtower auto-deploy) (push) Successful in 7m0s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 1m9s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 3m7s
2026-07-13 11:07:10 +08:00
CI Bot acca081149 feat(worker): 直通渲染 stream copy 优化 - 无重编码性能提升10倍+
CI/CD Pipeline / Frontend Lint (pull_request) Successful in 13m57s
CI/CD Pipeline / Validate Code Quality And Tests (pull_request) Failing after 18m2s
CI/CD Pipeline / Build & Push Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Build Production Runtime Images (pull_request) Has been skipped
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Integration Tests (pull_request) Successful in 10m58s
当输入片段与输出参数完全一致(h264/yuv420p/同分辨率/同帧率/无字幕/无特效)时,
走 FFmpeg stream copy 不重编码,性能提升 10 倍以上(典型场景 20s → 1-2s)。

核心改动:
1. probe_video_info 增强:返回编码/像素格式/音频信息(用于 copy 条件判断)
2. _can_use_stream_copy:7项条件检查(编码/分辨率/帧率/像素格式/字幕/trim等)
3. _try_render_stream_copy:stream copy 执行,失败自动回退到重编码
4. 单clip直通场景先尝试 stream copy,不满足或失败再回退带滤镜的直通渲染
5. trim 用 -ss/-t 实现,无需滤镜,copy 模式下也能用

安全保障:
- 条件不满足自动跳过,不影响现有渲染质量
- FFmpeg 失败自动回退到重编码,不影响成功率
- 输出文件为空时判定失败并清理损坏文件

测试:9个新增单测 + 现有75个统一渲染测试,全绿 
(含条件判断、成功路径、失败回退、完整渲染流程等场景)
2026-07-13 10:37:11 +08:00
xiaoxia bbe831f9e0 feat(worker): render_edit_plan 接入 Feature Flag 灰度控制 (#240) (#240)
CI/CD Pipeline / Frontend Lint (push) Successful in 10m57s
CI/CD Pipeline / Validate Code Quality And Tests (push) Successful in 13m35s
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 / Build & Push Staging (Watchtower auto-deploy) (push) Successful in 5m21s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 2m20s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 3m41s
CI/CD Pipeline / Integration Tests (push) Successful in 17m2s
2026-07-13 10:27:25 +08:00
CI Bot ffed853f6e fix(worker): FFmpeg 超时保护 - 防止渲染 hang 住导致 worker 永久阻塞
CI/CD Pipeline / Frontend Lint (pull_request) Successful in 10m46s
CI/CD Pipeline / Validate Code Quality And Tests (pull_request) Successful in 13m0s
CI/CD Pipeline / Build & Push Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Build Production Runtime Images (pull_request) Has been skipped
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Integration Tests (pull_request) Successful in 15m43s
根因:新引擎 UnifiedRenderService 所有 FFmpeg 调用通过 run_ffmpeg 执行,
但 subprocess.run 未设置 timeout,FFmpeg hang 住时 worker 线程永久阻塞。

旧引擎 compose_video 有单独的 timeout=3600,但新引擎路径没有。

修复:
1. run_ffmpeg 新增默认超时 1800s(30分钟),支持自定义传参
2. 捕获 TimeoutExpired 并打 error 日志后重新抛出
3. probe_video_info 新增 timeout=15s 超时保护
4. 7个单元测试覆盖超时逻辑

影响范围:所有通过 run_ffmpeg 调用的 FFmpeg 命令
(unified_render_service 所有渲染/混音/合并操作)
2026-07-13 10:17:14 +08:00
CI Bot f0cf4f50dc feat(api): 新增渲染结果内部下载接口 - 支持按视频ID/任务ID获取预签名URL
CI/CD Pipeline / Frontend Lint (pull_request) Successful in 12m19s
CI/CD Pipeline / Validate Code Quality And Tests (pull_request) Successful in 17m19s
CI/CD Pipeline / Build & Push Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Build Production Runtime Images (pull_request) Has been skipped
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Integration Tests (pull_request) Successful in 13m13s
- 新增 /internal/render/videos/{video_id}/download-url 接口:单个视频下载URL
- 新增 /internal/render/tasks/{task_id}/videos 接口:任务下所有视频+下载URL
- 内部 API Key 鉴权(X-API-Key),不对外开放
- 下载URL有效期24小时,方便对比工具批量下载
- 支持按 status 筛选(completed/failed 等)
- 8个单元测试覆盖核心逻辑
2026-07-13 09:54:18 +08:00
xiaoxia 4c5ab7f80e fix(worker): OSS上传崩溃修复 - connect超时 + 分片上传 + 总超时保护 (#238)
CI/CD Pipeline / Frontend Lint (push) Successful in 11m1s
CI/CD Pipeline / Validate Code Quality And Tests (push) Successful in 14m59s
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 / Build & Push Staging (Watchtower auto-deploy) (push) Successful in 5m35s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 2m21s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 3m53s
CI/CD Pipeline / Integration Tests (push) Successful in 18m41s
2026-07-13 09:33:13 +08:00
xiaoxia 9b2e782abd fix(worker): 修复 dedup fingerprint JSON 序列化失败 - np.float32 转原生 float (#239)
CI/CD Pipeline / Validate Code Quality And Tests (push) Successful in 1m19s
CI/CD Pipeline / Frontend Lint (push) Successful in 50s
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 1m14s
CI/CD Pipeline / Build & Push Staging (Watchtower auto-deploy) (push) Successful in 48m23s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 3m16s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 6m20s
2026-07-13 08:54:58 +08:00
12 changed files with 1505 additions and 128 deletions
+5
View File
@@ -13,6 +13,7 @@ from app.api.routes.generated_videos import router as generated_videos_router
from app.api.routes.generation_tasks import router as generation_tasks_router
from app.api.routes.health import router as health_check_router
from app.api.routes.ingest_jobs import router as ingest_jobs_router
from app.api.routes.internal_render import router as internal_render_router
from app.api.routes.jobs import router as jobs_router
from app.api.routes.projects import router as projects_router
from app.api.routes.recipes import router as recipes_router
@@ -156,3 +157,7 @@ api_router.include_router(
feature_flags_router,
tags=["Internal"],
)
api_router.include_router(
internal_render_router,
tags=["Internal"],
)
+120
View File
@@ -0,0 +1,120 @@
"""渲染结果内部下载接口。
通过内部 API Key 鉴权,为灰度对比工具等内部系统提供渲染结果下载能力。
API:
GET /api/v1/internal/render/videos/{video_id}/download-url - 获取单个视频下载URL
GET /api/v1/internal/render/tasks/{task_id}/videos - 获取任务下所有视频及下载URL
鉴权:X-API-Key header,走内部 API Key 验证
"""
from __future__ import annotations
import logging
from typing import Any
from app.api.routes.auth import _verify_internal_api_key
from app.core.storage import OSSStorageService, get_storage_service
from app.dependencies import get_generated_video_repository
from fastapi import APIRouter, Depends, HTTPException, Query
from pydantic import BaseModel
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/internal/render", tags=["Internal"])
class InternalRenderVideoItem(BaseModel):
"""内部渲染视频项。"""
video_id: str
generation_task_id: str
project_id: str
name: str
file_url: str
file_size: int | None = None
duration: float | None = None
width: int | None = None
height: int | None = None
fps: float | None = None
status: str
download_url: str
class InternalRenderTaskVideosResponse(BaseModel):
"""任务下所有渲染视频响应。"""
task_id: str
count: int
videos: list[InternalRenderVideoItem]
class InternalRenderDownloadUrlResponse(BaseModel):
"""单个视频下载URL响应。"""
video_id: str
download_url: str
def _video_to_item(video: Any, download_url: str) -> InternalRenderVideoItem:
"""将 GeneratedVideo 领域对象转为响应项。"""
return InternalRenderVideoItem(
video_id=video.id,
generation_task_id=video.generation_task_id,
project_id=video.project_id,
name=video.name,
file_url=video.file_url,
file_size=getattr(video, "file_size", None),
duration=getattr(video, "duration", None),
width=getattr(video, "width", None),
height=getattr(video, "height", None),
fps=getattr(video, "fps", None),
status=video.status,
download_url=download_url,
)
@router.get("/videos/{video_id}/download-url", response_model=InternalRenderDownloadUrlResponse)
def get_render_video_download_url(
video_id: str,
_: bool = Depends(_verify_internal_api_key),
generated_video_repository: Any = Depends(get_generated_video_repository),
storage_service: OSSStorageService = Depends(get_storage_service),
) -> InternalRenderDownloadUrlResponse:
"""获取单个渲染视频的下载URL(预签名)。"""
video = generated_video_repository.get(video_id)
if video is None:
raise HTTPException(status_code=404, detail=f"GeneratedVideo {video_id} not found")
download_url = storage_service.get_download_url(video.file_url, expires_seconds=86400)
logger.info("内部渲染下载URL生成: video_id=%s", video_id)
return InternalRenderDownloadUrlResponse(video_id=video_id, download_url=download_url)
@router.get("/tasks/{task_id}/videos", response_model=InternalRenderTaskVideosResponse)
def get_render_task_videos(
task_id: str,
status: str | None = Query(None, description="按状态筛选,如 completed/failed"),
_: bool = Depends(_verify_internal_api_key),
generated_video_repository: Any = Depends(get_generated_video_repository),
storage_service: OSSStorageService = Depends(get_storage_service),
) -> InternalRenderTaskVideosResponse:
"""获取生成任务下所有渲染视频及下载URL。"""
videos = generated_video_repository.list_by_generation_task(task_id)
# 状态筛选
if status:
videos = [v for v in videos if v.status == status]
items = []
for video in videos:
download_url = storage_service.get_download_url(video.file_url, expires_seconds=86400)
items.append(_video_to_item(video, download_url))
logger.info("内部渲染任务视频查询: task_id=%s count=%d", task_id, len(items))
return InternalRenderTaskVideosResponse(
task_id=task_id,
count=len(items),
videos=items,
)
+7 -3
View File
@@ -93,12 +93,16 @@ class VideoFingerprint:
resolution: tuple[int, int]
def to_dict(self) -> dict:
# 注意:color_histograms 里的值可能是 np.float32(来自 cv2.normalize),
# 直接存进 dict 后 SQLAlchemy JSON 序列化会报 "float32 is not JSON serializable"。
# 这里统一转成 Python 原生 float。
native_histograms = [[float(v) for v in hist] for hist in self.color_histograms]
return {
"md5": self.md5,
"keyframe_phashes": self.keyframe_phashes,
"color_histograms": self.color_histograms,
"duration": self.duration,
"resolution": list(self.resolution),
"color_histograms": native_histograms,
"duration": float(self.duration),
"resolution": [int(self.resolution[0]), int(self.resolution[1])],
}
+44 -10
View File
@@ -42,6 +42,10 @@ XFADE_TRANSITION_MAP: dict[str, str] = {
DEFAULT_TRANSITION_DURATION = 0.5
# FFmpeg 执行默认超时(秒),防止 FFmpeg hang 住导致 worker 永久阻塞
# 默认 30 分钟,足够处理大部分短视频渲染;超长视频可单独传参覆盖
DEFAULT_FFMPEG_TIMEOUT = 1800
# ── FFmpeg 执行 ───────────────────────────────────────────────────────────────
@@ -50,12 +54,14 @@ def run_ffmpeg(
command: list[str],
*,
capture_output: bool = True,
timeout: int | None = DEFAULT_FFMPEG_TIMEOUT,
) -> tuple[str, str]:
"""执行 FFmpeg 命令。
Args:
command: 完整的 ffmpeg 命令列表(含 "ffmpeg" 本身)
capture_output: 是否捕获 stdout/stderr
timeout: 超时时间(秒),默认 1800s(30分钟);None 表示不设超时(不推荐)
Returns:
(stdout, stderr) 元组
@@ -63,6 +69,7 @@ def run_ffmpeg(
Raises:
subprocess.CalledProcessError: 命令执行失败时抛出,
异常信息包含完整 stderr 以便排查。
subprocess.TimeoutExpired: 超时未完成时抛出,FFmpeg 进程会被 kill。
"""
try:
result = subprocess.run( # nosec B603
@@ -71,8 +78,16 @@ def run_ffmpeg(
stdout=subprocess.PIPE if capture_output else None,
stderr=subprocess.PIPE if capture_output else None,
text=True,
timeout=timeout,
)
return (result.stdout or "", result.stderr or "")
except subprocess.TimeoutExpired as e:
logger.error(
"FFmpeg 命令超时 (%ds): command=%s",
timeout or -1,
" ".join(str(c) for c in command[:20]),
)
raise
except subprocess.CalledProcessError as e:
# 把完整 stderr 打到日志,方便排查 exit code 183 等问题
stderr_text = (e.stderr or "").strip()
@@ -148,10 +163,14 @@ def probe_duration(local_path: str | Path) -> float:
def probe_video_info(video_path: str) -> dict[str, Any]:
"""获取视频信息(宽、高、时长、fps)。
"""获取视频信息(宽、高、时长、fps、编码、像素格式)。
Returns:
{"width": int, "height": int, "duration": float, "fps": float}
{
"width": int, "height": int, "duration": float, "fps": float,
"video_codec": str, "audio_codec": str, "pix_fmt": str,
"has_audio": bool,
}
失败时返回默认值。
"""
try:
@@ -160,10 +179,8 @@ def probe_video_info(video_path: str) -> dict[str, Any]:
FFPROBE_BIN,
"-v",
"error",
"-select_streams",
"v:0",
"-show_entries",
"stream=width,height,r_frame_rate,duration",
"stream=width,height,r_frame_rate,duration,codec_name,codec_type,pix_fmt",
"-show_entries",
"format=duration",
"-of",
@@ -174,19 +191,25 @@ def probe_video_info(video_path: str) -> dict[str, Any]:
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
text=True,
timeout=15,
)
import json
info = json.loads(result.stdout)
stream = info.get("streams", [{}])[0]
streams = info.get("streams", [])
fmt = info.get("format", {})
width = int(stream.get("width", DEFAULT_OUTPUT_WIDTH))
height = int(stream.get("height", DEFAULT_OUTPUT_HEIGHT))
video_stream = next((s for s in streams if s.get("codec_type") == "video"), {})
audio_stream = next((s for s in streams if s.get("codec_type") == "audio"), {})
width = int(video_stream.get("width", DEFAULT_OUTPUT_WIDTH))
height = int(video_stream.get("height", DEFAULT_OUTPUT_HEIGHT))
video_codec = video_stream.get("codec_name", "") or ""
pix_fmt = video_stream.get("pix_fmt", "") or ""
# 解析帧率
fps_str = stream.get("r_frame_rate", "25/1")
fps_str = video_stream.get("r_frame_rate", "25/1")
if "/" in fps_str:
num, den = fps_str.split("/")
fps = float(num) / float(den) if float(den) > 0 else DEFAULT_FPS
@@ -194,13 +217,20 @@ def probe_video_info(video_path: str) -> dict[str, Any]:
fps = float(fps_str) if fps_str else DEFAULT_FPS
# 时长
duration = float(fmt.get("duration", 0)) or float(stream.get("duration", 0))
duration = float(fmt.get("duration", 0)) or float(video_stream.get("duration", 0))
has_audio = bool(audio_stream)
audio_codec = audio_stream.get("codec_name", "") or ""
return {
"width": width,
"height": height,
"duration": duration,
"fps": round(fps, 2),
"video_codec": video_codec,
"audio_codec": audio_codec,
"pix_fmt": pix_fmt,
"has_audio": has_audio,
}
except Exception as e:
logger.warning("获取视频信息失败: %s, error: %s", video_path, e)
@@ -209,6 +239,10 @@ def probe_video_info(video_path: str) -> dict[str, Any]:
"height": DEFAULT_OUTPUT_HEIGHT,
"duration": 0.0,
"fps": DEFAULT_FPS,
"video_codec": "",
"audio_codec": "",
"pix_fmt": "",
"has_audio": True,
}
+82 -10
View File
@@ -9,6 +9,7 @@ from __future__ import annotations
import hashlib
import logging
import os
import threading
from pathlib import Path
from typing import Optional
from urllib.parse import urlparse
@@ -17,6 +18,13 @@ import oss2
logger = logging.getLogger(__name__)
# OSS 上传配置
OSS_CONNECT_TIMEOUT = 10 # 连接超时(秒),防止 TCP 握手挂死
OSS_UPLOAD_TOTAL_TIMEOUT = 300 # 单文件上传总超时(秒),防止网络慢时无限卡住
OSS_MULTIPART_THRESHOLD = 100 * 1024 * 1024 # 分片上传阈值:100MB 以上走分片
OSS_PART_SIZE = 8 * 1024 * 1024 # 分片大小:8MB
OSS_MULTIPART_NUM_THREADS = 3 # 分片上传并发数
# ── OSS 配置 ──────────────────────────────────────────────────────────────────
@@ -43,6 +51,9 @@ def oss_bucket() -> oss2.Bucket | None:
P0-2 修复:endpoint 不带 scheme 时自动补 https:// 前缀,
确保 sign_url 等依赖 scheme 的方法返回 HTTPS URL。
P0-staging 修复:增加 connect_timeout=10s,防止网络抖动时
TCP 握手阶段无限挂死,导致 worker 进程卡死。
Returns:
oss2.Bucket 实例,配置缺失时返回 None。
"""
@@ -53,7 +64,12 @@ def oss_bucket() -> oss2.Bucket | None:
# endpoint 无 scheme 时补 https://,与 API 端 storage.py 保持一致
if not endpoint.startswith(("http://", "https://")):
endpoint = f"https://{endpoint}"
return oss2.Bucket(oss2.Auth(access_key_id, access_key_secret), endpoint, bucket_name)
return oss2.Bucket(
oss2.Auth(access_key_id, access_key_secret),
endpoint,
bucket_name,
connect_timeout=OSS_CONNECT_TIMEOUT,
)
def normalize_storage_key(storage_key_or_url: str) -> str:
@@ -96,6 +112,9 @@ def download_asset(asset_storage_key: str, local_path: Path) -> bool:
def upload_to_oss(local_path: Path, storage_key: str) -> str | None:
"""上传文件到 OSS,返回公开 URL。
大文件(>100MB)自动走分片上传,降低内存峰值,减少 OOM 风险。
上传加总超时保护(默认 300s),防止网络异常时无限挂死。
Args:
local_path: 本地文件路径
storage_key: 目标存储键
@@ -106,18 +125,71 @@ def upload_to_oss(local_path: Path, storage_key: str) -> str | None:
bucket = oss_bucket()
if bucket is None:
return None
try:
bucket.put_object_from_file(storage_key, str(local_path))
settings = oss_settings()
if settings:
_, _, endpoint, bucket_name = settings
endpoint_clean = endpoint.replace("https://", "").replace("http://", "")
return f"https://{bucket_name}.{endpoint_clean}/{storage_key}"
result: dict = {"url": None, "error": None, "file_size": 0}
done = threading.Event()
def _do_upload():
try:
# 尝试获取文件大小,用于分片判断和日志;stat 失败时 fallback 走普通上传
try:
file_size = local_path.stat().st_size
result["file_size"] = file_size
use_multipart = file_size >= OSS_MULTIPART_THRESHOLD
except OSError:
use_multipart = False
file_size = 0
if use_multipart:
# 分片上传:降低内存峰值,每片 8MB,3 线程并发
logger.info(
"大文件分片上传: storage_key=%s, size=%.1fMB, part_size=%dMB, threads=%d",
storage_key[:80],
file_size / 1024 / 1024,
OSS_PART_SIZE // 1024 // 1024,
OSS_MULTIPART_NUM_THREADS,
)
oss2.resumable_upload(
bucket,
storage_key,
str(local_path),
multipart_threshold=OSS_MULTIPART_THRESHOLD,
part_size=OSS_PART_SIZE,
num_threads=OSS_MULTIPART_NUM_THREADS,
)
else:
bucket.put_object_from_file(storage_key, str(local_path))
# 构造返回 URL
settings = oss_settings()
if settings:
_, _, endpoint, bucket_name = settings
endpoint_clean = endpoint.replace("https://", "").replace("http://", "")
result["url"] = f"https://{bucket_name}.{endpoint_clean}/{storage_key}"
except Exception as e:
result["error"] = e
logger.exception("上传 OSS 失败: %s", storage_key)
finally:
done.set()
upload_thread = threading.Thread(target=_do_upload, daemon=True)
upload_thread.start()
finished = done.wait(timeout=OSS_UPLOAD_TOTAL_TIMEOUT)
if not finished:
logger.error(
"OSS 上传超时(%.0fs),强制中止: storage_key=%s, size=%.1fMB",
OSS_UPLOAD_TOTAL_TIMEOUT,
storage_key[:80],
result["file_size"] / 1024 / 1024 if result["file_size"] else 0,
)
return None
except Exception:
logger.exception("上传 OSS 失败: %s", storage_key)
if result["error"]:
return None
return result["url"]
def get_signed_download_url(storage_key_or_url: str, expires_seconds: int = 3600) -> str | None:
"""生成预签名下载 URL(用于私有 bucket 的 URL 校验或临时下载)。
@@ -22,7 +22,6 @@
from __future__ import annotations
import logging
import os
import subprocess
import time
from dataclasses import dataclass, field
@@ -296,7 +295,7 @@ WrapStyle: 2
Encoding: UTF-8
[V4+ Styles]
Format: Name, Fontname, Fontsize, PrimaryColour, SecondaryColour, OutlineColour, BackColour, Bold, Italic, Underline, StrikeOut, ScaleX, ScaleY, Spacing, Angle, BorderStyle, Outline, Shadow, Alignment, MarginL, MarginR, MarginV, Encoding
Format: Name, Fontname, Fontsize, PrimaryColour, SecondaryColour, OutlineColour, BackColour, Bold, Italic, Underline, StrikeOut, ScaleX, ScaleY, Spacing, Angle, BorderStyle, Outline, Shadow, Alignment, MarginL, MarginR, MarginV, Encoding # noqa: E501
{chr(10).join(styles)}
[Events]
@@ -460,12 +459,25 @@ class UnifiedRenderService:
is_pass_through = self._can_use_pass_through(layers)
pass_through_has_audio = False
used_stream_copy = False
if is_pass_through:
# 直通优化:单clip场景一次FFmpeg同时处理视频+音频,省去提取+合并两次调用
pass_through_has_audio = self._render_pass_through(
# 先尝试 stream copy 优化(无重编码,性能提升 10 倍+)
# 条件不满足或失败时回退到带滤镜的直通渲染
stream_copy_ok = self._try_render_stream_copy(
layers, output_path, ass_path=ass_path, video_duration=video_duration
)
if stream_copy_ok:
used_stream_copy = True
# stream copy 模式下,直接探测输出是否有音频
clip = layers[0].clips[0]
info = probe_video_info(str(clip.local_path))
pass_through_has_audio = info.get("has_audio", True)
else:
# 回退到带滤镜的直通渲染
pass_through_has_audio = self._render_pass_through(
layers, output_path, ass_path=ass_path, video_duration=video_duration
)
else:
filter_complex, input_args = self._build_filter_complex(layers, ass_path=ass_path)
self._execute_ffmpeg(filter_complex, input_args, video_only_path)
@@ -473,10 +485,11 @@ class UnifiedRenderService:
t_video_end = time.time()
video_render_ms = int((t_video_end - t_video_start) * 1000)
logger.info(
"[unified-render] video render done: plan_id=%s duration_ms=%d pass_through=%s",
"[unified-render] video render done: plan_id=%s duration_ms=%d pass_through=%s stream_copy=%s",
self.plan.id,
video_render_ms,
is_pass_through,
used_stream_copy,
)
# 6. 音频后处理混音(直通场景已合并处理,跳过)
@@ -619,6 +632,176 @@ class UnifiedRenderService:
return False
return True
def _can_use_stream_copy(
self,
clip: ResolvedClip,
*,
ass_path: Path | None = None,
video_duration: float = 0.0,
) -> tuple[bool, str]:
"""判断是否可以走 stream copy(流拷贝,不重编码)。
性能提升:10 倍以上(典型场景从 20s → 1-2s)。
条件:
1. 视频编码为 h264(输出目标也是 h264)
2. 像素格式为 yuv420p
3. 分辨率与输出一致(不需要 scale/crop)
4. 帧率与输出一致(误差 < 0.1fps
5. 无字幕叠加(字幕需要滤镜)
6. 无 trim 需求(或 trim 后恰好等于原时长)
7. 无转场、无特效(单 clip 直通已保证)
Returns:
(是否可以 copy, 原因说明)
"""
# 有字幕 → 需要滤镜 → 不能 copy
if ass_path is not None:
return False, "有字幕叠加"
# 探测输入视频参数
info = probe_video_info(str(clip.local_path))
# 编码必须是 h264
if info.get("video_codec", "") != "h264":
return False, f"视频编码不是h264: {info.get('video_codec', 'unknown')}"
# 像素格式必须是 yuv420p
if info.get("pix_fmt", "") != "yuv420p":
return False, f"像素格式不是yuv420p: {info.get('pix_fmt', 'unknown')}"
# 分辨率必须一致
if info.get("width", 0) != self.output_width or info.get("height", 0) != self.output_height:
return False, (
f"分辨率不匹配: "
f"{info.get('width', 0)}x{info.get('height', 0)} "
f"vs {self.output_width}x{self.output_height}"
)
# 帧率必须一致(误差 < 0.1fps
fps_diff = abs(info.get("fps", 0) - self.output_fps)
if fps_diff > 0.1:
return False, f"帧率不匹配: {info.get('fps', 0)} vs {self.output_fps}"
# 检查是否需要 trim
effective_duration = UnifiedRenderService._clip_effective_duration(clip)
if effective_duration > 0:
# 有 trim 需求但视频时长足够,可用 -ss/-t 实现 copy trim
input_duration = info.get("duration", 0)
if input_duration <= 0:
return False, "无法探测输入时长"
# trim 起始点 + 目标时长 <= 输入时长
start_time = getattr(clip, "start_time", 0) or 0
if start_time + effective_duration > input_duration + 0.1:
return False, "trim 超出输入时长"
# video_duration 截断
if video_duration > 0 and effective_duration > 0:
final_duration = min(effective_duration, video_duration)
if final_duration != effective_duration:
# 也需要截断,但 -t 可以 copy 模式下用
pass
return True, "所有条件满足"
def _try_render_stream_copy(
self,
layers: list[RenderLayer],
output_path: Path,
*,
ass_path: Path | None = None,
video_duration: float = 0.0,
) -> bool:
"""尝试 stream copy 渲染,成功返回 True,失败返回 False(调用方回退到重编码)。
stream copy 模式:不重编码,直接拷贝视频/音频流,性能提升 10 倍+。
仅用于单 clip 直通场景且满足 copy 条件。
"""
clip = layers[0].clips[0]
role = layers[0].role
# 判断是否满足 copy 条件
can_copy, reason = self._can_use_stream_copy(clip, ass_path=ass_path, video_duration=video_duration)
if not can_copy:
logger.info(
"[unified-render] stream_copy 跳过: plan_id=%s reason=%s",
self.plan.id,
reason,
)
return False
# 构建 copy 命令
command = [
FFMPEG_BIN,
"-y",
]
# trim 支持(-ss 放在 -i 前 = input seeking,速度更快但精度稍差;
# 放在 -i 后 = output seeking,精度高但慢)
# 这里用 output seeking 保证精度,反正 copy 模式已经很快了
start_time = getattr(clip, "start_time", 0) or 0
effective_duration = UnifiedRenderService._clip_effective_duration(clip)
command.extend(["-i", str(clip.local_path)])
if start_time > 0:
command.extend(["-ss", f"{start_time:.3f}"])
# 计算最终时长
final_duration = effective_duration
if video_duration > 0 and (final_duration <= 0 or final_duration > video_duration):
final_duration = video_duration
if final_duration > 0:
command.extend(["-t", f"{final_duration:.3f}"])
# 流拷贝
command.extend(
[
"-c:v",
"copy",
"-c:a",
"copy",
"-movflags",
"+faststart",
str(output_path),
]
)
logger.info(
"[unified-render] stream_copy 渲染: plan_id=%s clip=%s role=%s duration=%.2fs",
self.plan.id,
clip.clip_id,
role,
final_duration,
)
try:
run_ffmpeg(command)
# 验证输出文件存在且有大小
if output_path.exists() and output_path.stat().st_size > 0:
logger.info(
"[unified-render] stream_copy 成功: plan_id=%s size=%d",
self.plan.id,
output_path.stat().st_size,
)
return True
else:
logger.warning("[unified-render] stream_copy 输出为空: plan_id=%s", self.plan.id)
return False
except (subprocess.CalledProcessError, subprocess.TimeoutExpired) as e:
logger.warning(
"[unified-render] stream_copy 失败,回退到重编码: plan_id=%s error=%s",
self.plan.id,
str(e)[:200],
)
# 清理可能的损坏输出文件
if output_path.exists():
try:
output_path.unlink()
except OSError:
pass
return False
def _render_pass_through(
self,
layers: list[RenderLayer],
@@ -797,7 +980,6 @@ class UnifiedRenderService:
# 计算 PiP 位置
pip_width = int(self.output_width * _PIP_SCALE)
pip_height = int(self.output_height * _PIP_SCALE)
margin = 20 # 边距
if "overlay" in layer_map:
+302 -99
View File
@@ -1,13 +1,18 @@
"""剪辑计划渲染任务 — Phase 8 任务 2.05.
"""剪辑计划渲染任务 — 支持 Feature Flag 灰度.
Celery 任务 worker.render_edit_plan:
1. 加载 EditPlan + EditPlanClips
2. 下载各片段素材
3. 使用 UnifiedRenderService 按时间线+图层渲染
2. 根据 Feature Flag 选择渲染引擎(legacy / unified
3. 下载各片段素材 + 渲染
4. 上传渲染结果到 OSS
5. 创建 GeneratedVideo 记录 + 查重
6. 更新 EditPlan / EditPlanClip 状态
7. 更新 GenerationTask 进度
渲染引擎灰度:
- 走 Feature Flag (render_engine) 控制
- legacy: VideoComposeService + FFmpeg filter_complex
- unified: UnifiedRenderService 图层架构
"""
from __future__ import annotations
@@ -63,14 +68,268 @@ def _get_repos():
# ── Celery Task ───────────────────────────────────────────────────────────────
def _resolve_render_engine(user_id: str) -> str:
"""根据 Feature Flag 决定使用哪个渲染引擎。
Returns:
"legacy""unified"
"""
try:
from video_processing.render_engine_resolver import get_render_engine_resolver
resolver = get_render_engine_resolver()
return resolver.get_engine(user_id=user_id)
except Exception as exc:
logger.warning("获取渲染引擎配置失败,fallback 到 legacy: %s", exc)
return "legacy"
def _mark_plan_failed(plan_repo, plan_id, gen_task_repo, generation_task_id, error_msg: str):
"""统一的计划失败标记工具。"""
plan = plan_repo.get(plan_id)
if plan and plan.status.value == "rendering":
plan.mark_failed()
plan_repo.update(plan)
if generation_task_id:
gen_task = gen_task_repo.get(generation_task_id)
if gen_task and gen_task.status.value != "failed":
gen_task.status = "failed"
gen_task.error_message = error_msg
gen_task.completed_at = datetime.now(timezone.utc)
gen_task_repo.update(gen_task)
def _finalize_render_success(
plan,
plan_repo,
clip_repo,
gen_task_repo,
db,
plan_id: str,
output_url: str,
storage_key: str,
duration: float,
file_size: int,
width: int,
height: int,
rendered_clip_ids: list[str],
failed_clip_ids: list[str],
generation_task_id: str,
output_path: Path,
engine: str,
) -> dict:
"""渲染成功后的统一收尾:查重 + 更新状态 + 返回结果。"""
# 创建 GeneratedVideo 记录 + 查重
project_id = plan.project_id or ""
batch_id = plan.config.get("batch_id", "")
mode = plan.config.get("mode", "edit_plan")
if generation_task_id and project_id:
try:
create_video_record_and_dedup(
generation_task_id=generation_task_id,
project_id=project_id,
batch_id=batch_id,
file_url=output_url or "",
file_size=file_size,
duration=duration,
video_path=str(output_path),
mode=mode,
session=db,
width=width,
height=height,
fps=OUTPUT_FPS,
)
except Exception as dedup_err:
logger.warning("查重失败(不影响渲染结果): %s", dedup_err)
# 更新片段状态为 rendered
for clip_id in rendered_clip_ids:
clip = clip_repo.get(clip_id)
if clip and clip.status.value == "ready":
clip.mark_rendered()
clip_repo.update(clip)
# 更新 EditPlan 状态为 completed
plan.config["rendered_url"] = output_url or ""
plan.config["rendered_storage_key"] = storage_key
plan.mark_completed()
plan_repo.update(plan)
# 更新 GenerationTask 状态为 completed
if generation_task_id:
gen_task = gen_task_repo.get(generation_task_id)
if gen_task:
gen_task.status = "completed"
gen_task.progress = 100.0
gen_task.result_count = len(rendered_clip_ids)
gen_task.completed_at = datetime.now(timezone.utc)
gen_task_repo.update(gen_task)
logger.info(
"剪辑计划渲染完成: plan_id=%s engine=%s rendered=%d failed=%d duration=%.1fs",
plan_id,
engine,
len(rendered_clip_ids),
len(failed_clip_ids),
duration,
)
return {
"status": "completed",
"plan_id": plan_id,
"rendered_count": len(rendered_clip_ids),
"failed_count": len(failed_clip_ids),
"output_url": output_url,
"duration": duration,
}
def _render_with_unified(
plan,
clips,
asset_path_map: dict[str, Path],
tmpdir_path: Path,
rendered_clip_ids: list[str],
plan_id: str,
generation_task_id: str,
plan_repo,
clip_repo,
gen_task_repo,
db,
) -> dict:
"""统一渲染引擎路径(UnifiedRenderService 图层架构)。"""
render_service = UnifiedRenderService(
plan=plan,
clips=clips,
asset_path_map=asset_path_map,
work_dir=tmpdir_path,
output_width=OUTPUT_WIDTH,
output_height=OUTPUT_HEIGHT,
output_fps=int(OUTPUT_FPS),
)
try:
render_result = render_service.render()
except Exception as render_err:
logger.error("渲染失败(unified): %s%s", plan_id, render_err)
_mark_plan_failed(plan_repo, plan_id, gen_task_repo, generation_task_id, f"渲染失败: {render_err}")
return {"status": "error", "message": f"渲染失败: {render_err}"}
output_path = render_result.output_path
# 上传到 OSS
storage_key = f"rendered/{plan_id}/output.mp4"
output_url = upload_to_oss(output_path, storage_key)
failed_clip_ids: list[str] = []
return _finalize_render_success(
plan=plan,
plan_repo=plan_repo,
clip_repo=clip_repo,
gen_task_repo=gen_task_repo,
db=db,
plan_id=plan_id,
output_url=output_url or "",
storage_key=storage_key,
duration=render_result.duration,
file_size=render_result.file_size,
width=render_result.width,
height=render_result.height,
rendered_clip_ids=rendered_clip_ids,
failed_clip_ids=failed_clip_ids,
generation_task_id=generation_task_id,
output_path=output_path,
engine="unified",
)
def _render_with_legacy(
plan,
clips,
rendered_clip_ids: list[str],
failed_clip_ids: list[str],
tmpdir_path: Path,
plan_id: str,
generation_task_id: str,
plan_repo,
clip_repo,
gen_task_repo,
db,
) -> dict:
"""旧引擎路径(VideoComposeService + FFmpeg filter_complex)。"""
import os
import subprocess
from apps.api.app.services.video_compose_service import VideoComposeService
compose_svc = VideoComposeService(db)
# 校验合成条件
validation = compose_svc.validate_compose(plan_id)
if not validation.valid:
error_msg = "; ".join(validation.errors)
logger.error("合成校验失败(legacy): %s%s", plan_id, error_msg)
_mark_plan_failed(plan_repo, plan_id, gen_task_repo, generation_task_id, f"合成校验失败: {error_msg}")
return {"status": "error", "message": error_msg}
# 构建 FFmpeg 命令
output_dir = os.environ.get("VIDEO_OUTPUT_DIR", str(tmpdir_path))
output_path = Path(output_dir) / f"{plan_id}.mp4"
compose_cmd = compose_svc.build_compose_command(plan_id, str(output_path))
logger.info("执行 FFmpeg (legacy): plan_id=%s", plan_id)
try:
subprocess.run(
compose_cmd.command,
check=True,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
text=True,
timeout=3600,
)
except subprocess.CalledProcessError as e:
error_msg = f"FFmpeg 执行失败: {e.stderr[:500]}"
logger.error("FFmpeg 执行失败(legacy): %s%s", plan_id, error_msg)
_mark_plan_failed(plan_repo, plan_id, gen_task_repo, generation_task_id, error_msg)
return {"status": "error", "message": error_msg}
# 获取文件大小
file_size = output_path.stat().st_size if output_path.exists() else 0
duration = compose_cmd.estimated_duration or 0.0
# 上传到 OSS
storage_key = f"rendered/{plan_id}/output.mp4"
output_url = upload_to_oss(output_path, storage_key)
return _finalize_render_success(
plan=plan,
plan_repo=plan_repo,
clip_repo=clip_repo,
gen_task_repo=gen_task_repo,
db=db,
plan_id=plan_id,
output_url=output_url or "",
storage_key=storage_key,
duration=duration,
file_size=file_size,
width=OUTPUT_WIDTH,
height=OUTPUT_HEIGHT,
rendered_clip_ids=rendered_clip_ids,
failed_clip_ids=failed_clip_ids,
generation_task_id=generation_task_id,
output_path=output_path,
engine="legacy",
)
@celery_app.task(name="worker.render_edit_plan", bind=True, max_retries=2)
def render_edit_plan(self, plan_id: str) -> dict:
"""渲染剪辑计划
流程:
1. 加载 EditPlan + EditPlanClips
2. 下载各片段素材到临时目录,构建 asset_path_map
3. 使用 UnifiedRenderService 按时间线+图层渲染
2. 根据 Feature Flag 选择渲染引擎(legacy / unified
3. 下载素材 + 渲染
4. 上传渲染结果到 OSS
5. 创建 GeneratedVideo 记录 + 查重
6. 更新 EditPlan → completed, EditPlanClips → rendered
@@ -79,6 +338,7 @@ def render_edit_plan(self, plan_id: str) -> dict:
logger.info("开始渲染剪辑计划: plan_id=%s", plan_id)
generation_task_id = ""
engine = "legacy"
for repos in _get_repos():
plan_repo, clip_repo, gen_task_repo, db = repos
@@ -93,7 +353,12 @@ def render_edit_plan(self, plan_id: str) -> dict:
# 获取 generation_task_id(提前读取,确保 except 块可用)
generation_task_id = plan.config.get("generation_task_id", "")
# 2. 加载片段列表(按 order 排序
# 2. 选择渲染引擎(Feature Flag 灰度控制
user_id = plan.created_by_user_id or ""
engine = _resolve_render_engine(user_id)
logger.info("剪辑计划渲染引擎: plan_id=%s engine=%s user_id=%s", plan_id, engine, user_id)
# 3. 加载片段列表(按 order 排序)
clips = clip_repo.list_by_plan(plan_id, skip=0, limit=10000)
if not clips:
logger.warning("剪辑计划没有片段: %s", plan_id)
@@ -174,100 +439,38 @@ def render_edit_plan(self, plan_id: str) -> dict:
gen_task_repo.update(gen_task)
return {"status": "error", "message": "所有片段素材下载失败"}
# 4. 使用 UnifiedRenderService 渲染
render_service = UnifiedRenderService(
plan=plan,
clips=clips,
asset_path_map=asset_path_map,
work_dir=tmpdir_path,
output_width=OUTPUT_WIDTH,
output_height=OUTPUT_HEIGHT,
output_fps=int(OUTPUT_FPS),
)
# 4. 根据引擎选择渲染方式
if engine == "unified":
result = _render_with_unified(
plan=plan,
clips=clips,
asset_path_map=asset_path_map,
tmpdir_path=tmpdir_path,
rendered_clip_ids=rendered_clip_ids,
plan_id=plan_id,
generation_task_id=generation_task_id,
plan_repo=plan_repo,
clip_repo=clip_repo,
gen_task_repo=gen_task_repo,
db=db,
)
else:
result = _render_with_legacy(
plan=plan,
clips=clips,
rendered_clip_ids=rendered_clip_ids,
failed_clip_ids=failed_clip_ids,
tmpdir_path=tmpdir_path,
plan_id=plan_id,
generation_task_id=generation_task_id,
plan_repo=plan_repo,
clip_repo=clip_repo,
gen_task_repo=gen_task_repo,
db=db,
)
try:
render_result = render_service.render()
except Exception as render_err:
logger.error("渲染失败: %s%s", plan_id, render_err)
plan.mark_failed()
plan_repo.update(plan)
if generation_task_id:
gen_task = gen_task_repo.get(generation_task_id)
if gen_task:
gen_task.status = "failed"
gen_task.error_message = f"渲染失败: {render_err}"
gen_task.completed_at = datetime.now(timezone.utc)
gen_task_repo.update(gen_task)
return {"status": "error", "message": f"渲染失败: {render_err}"}
output_path = render_result.output_path
# 5. 上传到 OSS
storage_key = f"rendered/{plan_id}/output.mp4"
output_url = upload_to_oss(output_path, storage_key)
# 6. 创建 GeneratedVideo 记录 + 查重
project_id = plan.project_id or ""
batch_id = plan.config.get("batch_id", "")
mode = plan.config.get("mode", "edit_plan")
if generation_task_id and project_id:
try:
create_video_record_and_dedup(
generation_task_id=generation_task_id,
project_id=project_id,
batch_id=batch_id,
file_url=output_url or "",
file_size=render_result.file_size,
duration=render_result.duration,
video_path=str(output_path),
mode=mode,
session=db,
width=render_result.width,
height=render_result.height,
fps=OUTPUT_FPS,
)
except Exception as dedup_err:
logger.warning("查重失败(不影响渲染结果): %s", dedup_err)
# 7. 更新片段状态为 rendered
for clip_id in rendered_clip_ids:
clip = clip_repo.get(clip_id)
if clip and clip.status.value == "ready":
clip.mark_rendered()
clip_repo.update(clip)
# 8. 更新 EditPlan 状态为 completed
plan.config["rendered_url"] = output_url or ""
plan.config["rendered_storage_key"] = storage_key
plan.mark_completed()
plan_repo.update(plan)
# 9. 更新 GenerationTask 状态为 completed
if generation_task_id:
gen_task = gen_task_repo.get(generation_task_id)
if gen_task:
gen_task.status = "completed"
gen_task.progress = 100.0
gen_task.result_count = len(rendered_clip_ids)
gen_task.completed_at = datetime.now(timezone.utc)
gen_task_repo.update(gen_task)
logger.info(
"剪辑计划渲染完成: plan_id=%s rendered=%d failed=%d duration=%.1fs",
plan_id,
len(rendered_clip_ids),
len(failed_clip_ids),
render_result.duration,
)
return {
"status": "completed",
"plan_id": plan_id,
"rendered_count": len(rendered_clip_ids),
"failed_count": len(failed_clip_ids),
"output_url": output_url,
"duration": render_result.duration,
}
result["engine"] = engine
return result
except Exception as exc:
logger.exception("渲染剪辑计划异常: %s", plan_id)
+2
View File
@@ -63,6 +63,8 @@ class StubEditPlan:
template_id: str = "tmpl-001"
status: Any = None
config: dict = field(default_factory=dict)
project_id: str = ""
created_by_user_id: str = "user-001"
def mark_failed(self):
self.status = _StubStatus("failed")
@@ -0,0 +1,93 @@
"""FFmpeg 超时保护测试。
验证 run_ffmpeg / probe_video_info 的超时保护机制,
防止 FFmpeg hang 住导致 worker 永久阻塞。
"""
from __future__ import annotations
import subprocess
from unittest.mock import MagicMock, patch
import pytest
from video_processing.ffmpeg_utils import (
DEFAULT_FFMPEG_TIMEOUT,
probe_video_info,
run_ffmpeg,
)
# ── run_ffmpeg 超时保护 ──────────────────────────────────────────────────────
class TestRunFFmpegTimeout:
"""run_ffmpeg 超时保护测试。"""
def test_default_timeout_is_set(self):
"""默认超时应为 1800 秒(30分钟)。"""
assert DEFAULT_FFMPEG_TIMEOUT == 1800
def test_timeout_expired_is_raised(self):
"""超时未完成时 TimeoutExpired 异常被传播。"""
with patch("video_processing.ffmpeg_utils.subprocess.run") as mock_run:
mock_run.side_effect = subprocess.TimeoutExpired(cmd=["ffmpeg", "test"], timeout=1)
with pytest.raises(subprocess.TimeoutExpired):
run_ffmpeg(["ffmpeg", "test"])
def test_custom_timeout(self):
"""支持自定义超时时间。"""
with patch("video_processing.ffmpeg_utils.subprocess.run") as mock_run:
mock_run.side_effect = subprocess.TimeoutExpired(cmd=["ffmpeg"], timeout=5)
with pytest.raises(subprocess.TimeoutExpired):
run_ffmpeg(["ffmpeg", "test"], timeout=5)
def test_none_timeout_disables_protection(self):
"""timeout=None 可以禁用超时保护(不推荐)。"""
with patch("video_processing.ffmpeg_utils.subprocess.run") as mock_run:
mock_result = MagicMock()
mock_result.stdout = ""
mock_result.stderr = ""
mock_run.return_value = mock_result
run_ffmpeg(["ffmpeg", "test"], timeout=None)
# 验证 timeout=None 被传递
call_kwargs = mock_run.call_args.kwargs
assert call_kwargs["timeout"] is None
def test_called_process_error_still_raised(self):
"""超时异常不影响原有 CalledProcessError 的抛出。"""
with patch("video_processing.ffmpeg_utils.subprocess.run") as mock_run:
mock_run.side_effect = subprocess.CalledProcessError(returncode=1, cmd=["ffmpeg"], stderr="error msg")
with pytest.raises(subprocess.CalledProcessError):
run_ffmpeg(["ffmpeg", "test"])
# ── probe_video_info 超时保护 ────────────────────────────────────────────────
class TestProbeVideoInfoTimeout:
"""probe_video_info 超时保护测试。"""
def test_probe_uses_timeout(self):
"""probe_video_info 调用 ffprobe 时应设置 timeout=15。"""
with patch("video_processing.ffmpeg_utils.subprocess.run") as mock_run:
mock_run.side_effect = subprocess.TimeoutExpired(cmd=["ffprobe"], timeout=15)
# 超时异常被捕获,返回默认值
result = probe_video_info("/tmp/test.mp4")
assert result["width"] == 1280 # DEFAULT_OUTPUT_WIDTH
assert result["height"] == 720 # DEFAULT_OUTPUT_HEIGHT
def test_probe_success(self):
"""正常情况应解析 ffprobe JSON 输出。"""
fake_output = """
{
"streams": [{"width": 1920, "height": 1080, "r_frame_rate": "30/1", "duration": "10.5"}],
"format": {"duration": "10.5"}
}
"""
with patch("video_processing.ffmpeg_utils.subprocess.run") as mock_run:
mock_result = MagicMock()
mock_result.stdout = fake_output
mock_run.return_value = mock_result
result = probe_video_info("/tmp/test.mp4")
assert result["width"] == 1920
assert result["height"] == 1080
assert abs(result["duration"] - 10.5) < 0.01
+190
View File
@@ -0,0 +1,190 @@
"""渲染结果内部下载接口单元测试。
测试 internal_render 路由的核心逻辑,mock 掉 repository 和 storage 依赖。
"""
from __future__ import annotations
from unittest.mock import MagicMock
import pytest
from app.api.routes.internal_render import (
InternalRenderDownloadUrlResponse,
InternalRenderTaskVideosResponse,
_video_to_item,
get_render_task_videos,
get_render_video_download_url,
)
# ── Helpers ────────────────────────────────────────────────────────────────
class MockVideo:
"""模拟 GeneratedVideo 领域对象。"""
def __init__(self, **kwargs):
self.id = kwargs.get("id", "video-1")
self.generation_task_id = kwargs.get("generation_task_id", "task-1")
self.project_id = kwargs.get("project_id", "proj-1")
self.name = kwargs.get("name", "test_video.mp4")
self.file_url = kwargs.get("file_url", "videos/test/output.mp4")
self.file_size = kwargs.get("file_size", 1024000)
self.duration = kwargs.get("duration", 30.5)
self.width = kwargs.get("width", 1080)
self.height = kwargs.get("height", 1920)
self.fps = kwargs.get("fps", 30.0)
self.status = kwargs.get("status", "completed")
# ── _video_to_item 测试 ────────────────────────────────────────────────────
class TestVideoToItem:
"""测试视频对象转响应项。"""
def test_basic_conversion(self):
video = MockVideo(id="v1", generation_task_id="t1", status="completed")
item = _video_to_item(video, "https://oss.example.com/download?v1")
assert item.video_id == "v1"
assert item.generation_task_id == "t1"
assert item.status == "completed"
assert item.download_url == "https://oss.example.com/download?v1"
def test_missing_optional_fields(self):
"""缺可选字段时返回 None。"""
video = MockVideo()
# 去掉可选字段
del video.file_size
del video.duration
item = _video_to_item(video, "https://example.com/dl")
assert item.file_size is None
assert item.duration is None
assert item.width == 1080 # 还在
# ── 路由函数测试 ────────────────────────────────────────────────────────────
class TestGetRenderVideoDownloadUrl:
"""测试单个视频下载URL接口。"""
def test_video_exists(self):
video = MockVideo(id="v-abc", file_url="videos/abc/out.mp4")
mock_repo = MagicMock()
mock_repo.get.return_value = video
mock_storage = MagicMock()
mock_storage.get_download_url.return_value = "https://oss.test/signed?v=abc"
result = get_render_video_download_url(
video_id="v-abc",
_=True,
generated_video_repository=mock_repo,
storage_service=mock_storage,
)
assert isinstance(result, InternalRenderDownloadUrlResponse)
assert result.video_id == "v-abc"
assert result.download_url == "https://oss.test/signed?v=abc"
mock_repo.get.assert_called_once_with("v-abc")
mock_storage.get_download_url.assert_called_once()
def test_video_not_found_raises_404(self):
from fastapi import HTTPException
mock_repo = MagicMock()
mock_repo.get.return_value = None
mock_storage = MagicMock()
with pytest.raises(HTTPException) as exc_info:
get_render_video_download_url(
video_id="nonexistent",
_=True,
generated_video_repository=mock_repo,
storage_service=mock_storage,
)
assert exc_info.value.status_code == 404
def test_download_url_long_expiry(self):
"""过期时间应为 24 小时(86400s)。"""
video = MockVideo(id="v1")
mock_repo = MagicMock()
mock_repo.get.return_value = video
mock_storage = MagicMock()
mock_storage.get_download_url.return_value = "https://oss.test/signed"
get_render_video_download_url(
video_id="v1",
_=True,
generated_video_repository=mock_repo,
storage_service=mock_storage,
)
# 验证 expires_seconds=86400
call_kwargs = mock_storage.get_download_url.call_args
assert call_kwargs.kwargs.get("expires_seconds") == 86400 or call_kwargs[1].get("expires_seconds") == 86400
class TestGetRenderTaskVideos:
"""测试任务视频列表接口。"""
def test_list_multiple_videos(self):
videos = [
MockVideo(id="v1", status="completed"),
MockVideo(id="v2", status="completed"),
MockVideo(id="v3", status="failed"),
]
mock_repo = MagicMock()
mock_repo.list_by_generation_task.return_value = videos
mock_storage = MagicMock()
mock_storage.get_download_url.return_value = "https://oss.test/signed"
result = get_render_task_videos(
task_id="task-1",
status=None,
_=True,
generated_video_repository=mock_repo,
storage_service=mock_storage,
)
assert isinstance(result, InternalRenderTaskVideosResponse)
assert result.task_id == "task-1"
assert result.count == 3
assert len(result.videos) == 3
def test_filter_by_status(self):
videos = [
MockVideo(id="v1", status="completed"),
MockVideo(id="v2", status="completed"),
MockVideo(id="v3", status="failed"),
]
mock_repo = MagicMock()
mock_repo.list_by_generation_task.return_value = videos
mock_storage = MagicMock()
mock_storage.get_download_url.return_value = "https://oss.test/signed"
result = get_render_task_videos(
task_id="task-1",
status="completed",
_=True,
generated_video_repository=mock_repo,
storage_service=mock_storage,
)
assert result.count == 2
assert all(v.status == "completed" for v in result.videos)
def test_empty_task(self):
mock_repo = MagicMock()
mock_repo.list_by_generation_task.return_value = []
mock_storage = MagicMock()
result = get_render_task_videos(
task_id="empty-task",
status=None,
_=True,
generated_video_repository=mock_repo,
storage_service=mock_storage,
)
assert result.count == 0
assert result.videos == []
+236
View File
@@ -0,0 +1,236 @@
"""P0-stagingOSS 上传崩溃修复测试.
测试:
1. oss_bucket() 传递 connect_timeout 参数
2. upload_to_oss() 小文件走 put_object_from_file,大文件走分片上传
3. upload_to_oss() 超时保护(超过总超时返回 None)
4. upload_to_oss() 异常时返回 None
"""
from __future__ import annotations
import os
import tempfile
import time
from pathlib import Path
from unittest.mock import MagicMock, patch
import pytest
# ── oss_bucket connect_timeout 测试 ───────────────────────────────────────────
class TestOSSBucketConnectTimeout:
"""测试 oss_bucket() 传递 connect_timeout 参数."""
def test_oss_bucket_has_connect_timeout(self):
"""oss_bucket 应传递 connect_timeout=10s 参数."""
from video_processing.oss_helpers import oss_bucket
mock_bucket_instance = MagicMock()
with (
patch.dict(
os.environ,
{
"OSS_ACCESS_KEY_ID": "test-key",
"OSS_ACCESS_KEY_SECRET": "test-secret",
"OSS_ENDPOINT": "oss-cn-hangzhou.aliyuncs.com",
"OSS_BUCKET_NAME": "test-bucket",
},
),
patch("video_processing.oss_helpers.oss2.Auth"),
patch("video_processing.oss_helpers.oss2.Bucket", return_value=mock_bucket_instance) as mock_bucket_cls,
):
bucket = oss_bucket()
assert bucket is mock_bucket_instance
# 验证 connect_timeout 关键字参数
call_kwargs = mock_bucket_cls.call_args[1]
assert "connect_timeout" in call_kwargs, "oss_bucket 应传递 connect_timeout 参数"
assert (
call_kwargs["connect_timeout"] == 10
), f"connect_timeout 应为 10,实际为 {call_kwargs['connect_timeout']}"
def test_oss_bucket_no_config_returns_none(self):
"""OSS 配置缺失时返回 None."""
from video_processing.oss_helpers import oss_bucket
with patch.dict(os.environ, {}, clear=True):
bucket = oss_bucket()
assert bucket is None
# ── upload_to_oss 分片上传测试 ────────────────────────────────────────────────
class TestUploadToOSSMultipart:
"""测试 upload_to_oss() 根据文件大小选择上传方式."""
def _create_temp_file(self, size_bytes: int) -> Path:
"""创建指定大小的临时文件."""
tmp = tempfile.NamedTemporaryFile(delete=False, suffix=".mp4")
tmp.write(b"x" * size_bytes)
tmp.close()
return Path(tmp.name)
def test_small_file_uses_put_object(self):
"""小文件(<100MB)走 put_object_from_file."""
from video_processing.oss_helpers import upload_to_oss
small_file = self._create_temp_file(10 * 1024 * 1024) # 10MB
try:
mock_bucket = MagicMock()
with (
patch.dict(
os.environ,
{
"OSS_ACCESS_KEY_ID": "test-key",
"OSS_ACCESS_KEY_SECRET": "test-secret",
"OSS_ENDPOINT": "oss-cn-hangzhou.aliyuncs.com",
"OSS_BUCKET_NAME": "test-bucket",
},
),
patch("video_processing.oss_helpers.oss2.Auth"),
patch("video_processing.oss_helpers.oss2.Bucket", return_value=mock_bucket),
patch("video_processing.oss_helpers.oss2.resumable_upload") as mock_resumable,
):
url = upload_to_oss(small_file, "test/small.mp4")
# 验证调用了 put_object_from_file
mock_bucket.put_object_from_file.assert_called_once()
# 验证没调用分片上传
mock_resumable.assert_not_called()
# 验证返回 URL
assert url == "https://test-bucket.oss-cn-hangzhou.aliyuncs.com/test/small.mp4"
finally:
small_file.unlink()
def test_large_file_uses_resumable_upload(self):
"""大文件(>=100MB)走 resumable_upload 分片上传."""
from video_processing.oss_helpers import upload_to_oss
large_file = self._create_temp_file(100 * 1024 * 1024) # 100MB
try:
mock_bucket = MagicMock()
with (
patch.dict(
os.environ,
{
"OSS_ACCESS_KEY_ID": "test-key",
"OSS_ACCESS_KEY_SECRET": "test-secret",
"OSS_ENDPOINT": "oss-cn-hangzhou.aliyuncs.com",
"OSS_BUCKET_NAME": "test-bucket",
},
),
patch("video_processing.oss_helpers.oss2.Auth"),
patch("video_processing.oss_helpers.oss2.Bucket", return_value=mock_bucket),
patch("video_processing.oss_helpers.oss2.resumable_upload") as mock_resumable,
):
url = upload_to_oss(large_file, "test/large.mp4")
# 验证调用了分片上传
mock_resumable.assert_called_once()
# 验证没调用 put_object_from_file
mock_bucket.put_object_from_file.assert_not_called()
# 验证分片参数
call_kwargs = mock_resumable.call_args[1]
assert call_kwargs["multipart_threshold"] == 100 * 1024 * 1024
assert call_kwargs["part_size"] == 8 * 1024 * 1024
assert call_kwargs["num_threads"] == 3
# 验证返回 URL
assert url == "https://test-bucket.oss-cn-hangzhou.aliyuncs.com/test/large.mp4"
finally:
large_file.unlink()
# ── upload_to_oss 超时测试 ────────────────────────────────────────────────────
class TestUploadToOSSTimeout:
"""测试 upload_to_oss() 超时保护."""
def test_upload_timeout_returns_none(self):
"""上传超过总超时时返回 None."""
from video_processing.oss_helpers import upload_to_oss
small_file = tempfile.NamedTemporaryFile(delete=False, suffix=".mp4")
small_file.write(b"x" * 1024) # 1KB
small_file.close()
file_path = Path(small_file.name)
def slow_upload(*args, **kwargs):
"""模拟慢速上传,超过超时时间."""
time.sleep(2)
mock_bucket = MagicMock()
mock_bucket.put_object_from_file.side_effect = slow_upload
try:
with (
patch.dict(
os.environ,
{
"OSS_ACCESS_KEY_ID": "test-key",
"OSS_ACCESS_KEY_SECRET": "test-secret",
"OSS_ENDPOINT": "oss-cn-hangzhou.aliyuncs.com",
"OSS_BUCKET_NAME": "test-bucket",
},
),
patch("video_processing.oss_helpers.oss2.Auth"),
patch("video_processing.oss_helpers.oss2.Bucket", return_value=mock_bucket),
patch("video_processing.oss_helpers.OSS_UPLOAD_TOTAL_TIMEOUT", 1), # 1秒超时
):
url = upload_to_oss(file_path, "test/slow.mp4")
# 超时应返回 None
assert url is None, "上传超时应返回 None"
finally:
file_path.unlink()
def test_upload_exception_returns_none(self):
"""上传异常时返回 None."""
from video_processing.oss_helpers import upload_to_oss
small_file = tempfile.NamedTemporaryFile(delete=False, suffix=".mp4")
small_file.write(b"x" * 1024)
small_file.close()
file_path = Path(small_file.name)
mock_bucket = MagicMock()
mock_bucket.put_object_from_file.side_effect = RuntimeError("Network error")
try:
with (
patch.dict(
os.environ,
{
"OSS_ACCESS_KEY_ID": "test-key",
"OSS_ACCESS_KEY_SECRET": "test-secret",
"OSS_ENDPOINT": "oss-cn-hangzhou.aliyuncs.com",
"OSS_BUCKET_NAME": "test-bucket",
},
),
patch("video_processing.oss_helpers.oss2.Auth"),
patch("video_processing.oss_helpers.oss2.Bucket", return_value=mock_bucket),
):
url = upload_to_oss(file_path, "test/error.mp4")
assert url is None, "上传异常应返回 None"
finally:
file_path.unlink()
def test_upload_no_bucket_returns_none(self):
"""OSS 未配置时返回 None."""
from video_processing.oss_helpers import upload_to_oss
small_file = tempfile.NamedTemporaryFile(delete=False, suffix=".mp4")
small_file.write(b"x" * 1024)
small_file.close()
file_path = Path(small_file.name)
try:
with patch.dict(os.environ, {}, clear=True):
url = upload_to_oss(file_path, "test/noconfig.mp4")
assert url is None
finally:
file_path.unlink()
+236
View File
@@ -1399,3 +1399,239 @@ class TestAudioMixing:
assert r1 is True and r2 is True and r3 is True
# 实际只探测了 1 次
assert mock_probe.call_count == 1
# ── 测试 stream copy 流拷贝优化 ───────────────────────────────────────────────
class TestStreamCopy:
"""stream copy 流拷贝优化测试。"""
def _make_single_clip_service(self):
clips = [_make_clip("c1", "main", order=0, duration=5.0)]
svc = _make_service(clips)
with _patch_path_exists(), patch("video_processing.unified_render_service.probe_duration", return_value=5.0):
resolved = svc._resolve_clips()
layers = svc._group_clips_into_layers(resolved)
return svc, resolved[0], layers
def test_can_use_stream_copy_all_conditions_met(self):
"""所有条件满足 → 可以 stream copy。"""
svc, clip, layers = self._make_single_clip_service()
probe_result = {
"width": 1280,
"height": 720,
"fps": 25.0,
"video_codec": "h264",
"pix_fmt": "yuv420p",
"duration": 5.0,
"has_audio": True,
"audio_codec": "aac",
}
with patch(
"video_processing.unified_render_service.probe_video_info",
return_value=probe_result,
):
can_copy, reason = svc._can_use_stream_copy(clip, ass_path=None, video_duration=0)
assert can_copy is True
assert "所有条件满足" in reason
def test_cannot_copy_with_subtitles(self):
"""有字幕 → 不能 stream copy。"""
svc, clip, layers = self._make_single_clip_service()
can_copy, reason = svc._can_use_stream_copy(clip, ass_path=Path("/tmp/sub.ass"), video_duration=0)
assert can_copy is False
assert "字幕" in reason
def test_cannot_copy_wrong_codec(self):
"""编码不是 h264 → 不能 stream copy。"""
svc, clip, layers = self._make_single_clip_service()
probe_result = {
"width": 1280,
"height": 720,
"fps": 25.0,
"video_codec": "hevc",
"pix_fmt": "yuv420p",
"duration": 5.0,
"has_audio": True,
"audio_codec": "aac",
}
with patch(
"video_processing.unified_render_service.probe_video_info",
return_value=probe_result,
):
can_copy, reason = svc._can_use_stream_copy(clip, ass_path=None, video_duration=0)
assert can_copy is False
assert "编码" in reason
def test_cannot_copy_wrong_resolution(self):
"""分辨率不匹配 → 不能 stream copy。"""
svc, clip, layers = self._make_single_clip_service()
probe_result = {
"width": 1920,
"height": 1080,
"fps": 25.0,
"video_codec": "h264",
"pix_fmt": "yuv420p",
"duration": 5.0,
"has_audio": True,
"audio_codec": "aac",
}
with patch(
"video_processing.unified_render_service.probe_video_info",
return_value=probe_result,
):
can_copy, reason = svc._can_use_stream_copy(clip, ass_path=None, video_duration=0)
assert can_copy is False
assert "分辨率" in reason
def test_cannot_copy_wrong_fps(self):
"""帧率不匹配 → 不能 stream copy。"""
svc, clip, layers = self._make_single_clip_service()
probe_result = {
"width": 1280,
"height": 720,
"fps": 30.0,
"video_codec": "h264",
"pix_fmt": "yuv420p",
"duration": 5.0,
"has_audio": True,
"audio_codec": "aac",
}
with patch(
"video_processing.unified_render_service.probe_video_info",
return_value=probe_result,
):
can_copy, reason = svc._can_use_stream_copy(clip, ass_path=None, video_duration=0)
assert can_copy is False
assert "帧率" in reason
def test_cannot_copy_wrong_pix_fmt(self):
"""像素格式不匹配 → 不能 stream copy。"""
svc, clip, layers = self._make_single_clip_service()
probe_result = {
"width": 1280,
"height": 720,
"fps": 25.0,
"video_codec": "h264",
"pix_fmt": "yuv422p",
"duration": 5.0,
"has_audio": True,
"audio_codec": "aac",
}
with patch(
"video_processing.unified_render_service.probe_video_info",
return_value=probe_result,
):
can_copy, reason = svc._can_use_stream_copy(clip, ass_path=None, video_duration=0)
assert can_copy is False
assert "像素格式" in reason
def test_try_render_stream_copy_success(self):
"""stream copy 渲染成功 → 返回 True。"""
svc, clip, layers = self._make_single_clip_service()
probe_result = {
"width": 1280,
"height": 720,
"fps": 25.0,
"video_codec": "h264",
"pix_fmt": "yuv420p",
"duration": 5.0,
"has_audio": True,
"audio_codec": "aac",
}
output_path = Path("/tmp/test_output.mp4")
def fake_stat():
m = MagicMock()
m.st_size = 1024000
return m
with (
patch(
"video_processing.unified_render_service.probe_video_info",
return_value=probe_result,
),
patch("video_processing.unified_render_service.run_ffmpeg") as mock_run,
patch("pathlib.Path.exists", return_value=True),
patch("pathlib.Path.stat", side_effect=fake_stat),
):
result = svc._try_render_stream_copy(layers, output_path, ass_path=None, video_duration=0)
assert result is True
mock_run.assert_called_once()
cmd = mock_run.call_args[0][0]
assert "-c:v" in cmd
assert "copy" in cmd
assert "-c:a" in cmd
def test_try_render_stream_copy_fallback_on_ffmpeg_error(self):
"""stream copy FFmpeg 失败 → 返回 False(调用方回退到重编码)。"""
svc, clip, layers = self._make_single_clip_service()
probe_result = {
"width": 1280,
"height": 720,
"fps": 25.0,
"video_codec": "h264",
"pix_fmt": "yuv420p",
"duration": 5.0,
"has_audio": True,
"audio_codec": "aac",
}
output_path = Path("/tmp/test_output.mp4")
import subprocess as sp
with (
patch(
"video_processing.unified_render_service.probe_video_info",
return_value=probe_result,
),
patch(
"video_processing.unified_render_service.run_ffmpeg",
side_effect=sp.CalledProcessError(1, ["ffmpeg"], stderr="copy failed"),
),
patch("pathlib.Path.exists", return_value=False),
):
result = svc._try_render_stream_copy(layers, output_path, ass_path=None, video_duration=0)
assert result is False
def test_render_uses_stream_copy_when_eligible(self):
"""完整渲染流程:满足条件时走 stream copy。"""
clips = [_make_clip("c1", "main", order=0, duration=5.0)]
svc = _make_service(clips)
probe_result = {
"width": 1280,
"height": 720,
"fps": 25.0,
"video_codec": "h264",
"pix_fmt": "yuv420p",
"duration": 5.0,
"has_audio": True,
"audio_codec": "aac",
}
def fake_stat():
m = MagicMock()
m.st_size = 1024000
return m
with (
_patch_path_exists(),
patch("video_processing.unified_render_service.probe_duration", return_value=5.0),
patch(
"video_processing.unified_render_service.probe_video_info",
return_value=probe_result,
),
patch("video_processing.unified_render_service.run_ffmpeg") as mock_run,
patch("pathlib.Path.stat", side_effect=fake_stat),
patch("shutil.copy2"),
):
result = svc.render()
assert mock_run.call_count == 1
cmd = mock_run.call_args[0][0]
assert "copy" in cmd
assert isinstance(result.output_path, Path)