feat(worker): 全GPU直连渲染,取消mezzanine CPU中间编码 #2086

Merged
xiaoxia merged 3 commits from feature/gpu-direct-pipeline into develop 2026-09-28 14:20:49 +08:00
5 changed files with 606 additions and 12 deletions
@@ -0,0 +1,296 @@
"""全 GPU 直连渲染管线(P1)。
背景:旧链路 worker 先用 CPU libx264 把 filter_complex 输出成 mezzanine(1080p 约 85s),
上传后再由 P4000 NVENC 编码,渲染后还要单独跑一次随机边缘裁剪重编码(约 26s)。
本管线取消 mezzanine:把原始素材签名 URL 作为多输入直接交给 P4000,filter_complex 内
一步完成 trim/scale/pad/concat/边缘裁剪/drawtext 字幕,末端 h264_nvenc 只编码一次;
TTS/BGM 音频也在同一命令里 amix 合成。
约束(P1):
- 仅覆盖智能剪辑主流场景:单一主视频轨、全硬切、无 PiP/overlay/水印/贴纸/片头片尾/绿幕。
不满足条件时调用方回退到现有 mezzanine/CPU 链路(功能不回归)。
- 字幕先用 drawtext(P4000 装好中文字体后可再切 subtitles 滤镜烧 ASS)。
"""
from __future__ import annotations
import logging
import uuid
from pathlib import Path
from typing import Any, Optional
logger = logging.getLogger(__name__)
# drawtext 默认字体名(fontconfig 解析);P4000 装好中文字体后可在 config 指定
DEFAULT_DRAWTEXT_FONT = "Noto Sans CJK SC"
def escape_drawtext_text(text: str) -> str:
"""转义 drawtext text= 中的特殊字符(ffmpeg 过滤器语法)。"""
if not text:
return ""
# 顺序重要:先转义反斜杠本身
s = text.replace("\\", "\\\\")
s = s.replace(":", "\\:")
s = s.replace("'", "\\'")
s = s.replace("%", "\\%")
s = s.replace(",", "\\,")
s = s.replace("[", "\\[").replace("]", "\\]")
s = s.replace(";", "\\;")
# 换行保留(drawtext 支持 %{...};真实换行需转成字面)
s = s.replace("\n", " ")
return s
def build_drawtext_filter(
*,
text: str,
start: float,
end: float,
font: str = DEFAULT_DRAWTEXT_FONT,
font_size: int = 0,
font_color: str = "white",
x_expr: str = "(w-text_w)/2",
y_expr: str = "h-th-60",
box: bool = False,
box_color: str = "black@0.5",
borderw: int = 0,
border_color: str = "black",
enable: bool = True,
) -> str:
"""构造单个 drawtext 滤镜字符串(不含输入/输出标签)。
[start, end] 秒的显示窗口通过 enable='between(t,...)' 控制。
"""
txt = escape_drawtext_text(text)
parts = [f"font={font}", f"text='{txt}'"]
if font_size and font_size > 0:
parts.append(f"fontsize={int(font_size)}")
parts.append(f"fontcolor={font_color}")
if box:
parts.append("box=1")
parts.append(f"boxcolor={box_color}")
if borderw and borderw > 0:
parts.append(f"borderw={int(borderw)}")
parts.append(f"bordercolor={border_color}")
parts.append(f"x={x_expr}")
parts.append(f"y={y_expr}")
if enable:
parts.append(f"enable='between(t,{start:.3f},{end:.3f})'")
return "drawtext=" + ":".join(parts)
def upload_local_audio_and_sign(
local_audio: Path,
*,
tmp_prefix: str = "tmp/gpu-direct-audio/",
expires: int = 3600,
) -> tuple[str, str]:
"""把本地音频(TTS/BGM)上传 OSS tmp 目录并签公网 GET URL。
Returns:
(signed_get_url, oss_key)
"""
from video_processing.oss_helpers import _storage # 类型: ignore
storage = _storage()
key = f"{tmp_prefix.rstrip('/')}/{uuid.uuid4().hex}{local_audio.suffix or '.mp3'}"
content_type = "audio/mpeg" if local_audio.suffix.lower() in (".mp3", ".mpeg") else "audio/mp4"
storage.upload_file(local_audio, key, content_type=content_type)
url = storage.get_download_url(key, expires)
return url, key
def sign_asset_url(storage_key: str, *, expires: int = 3600) -> str:
"""给原始素材 storage_key 签公网 GET URL(供 P4000 直接下载)。"""
from video_processing.oss_helpers import _storage # 类型: ignore
storage = _storage()
return storage.get_download_url(storage_key, expires)
# ── 直连渲染编排器 ────────────────────────────────────────────────────────
class DirectRenderPlan:
"""一次直连渲染的产物:inputs(裸文件名→URL)与完整 ffmpeg_args。"""
def __init__(self, inputs: dict[str, str], ffmpeg_args: list[str], oss_keys: list[str]):
self.inputs = inputs
self.ffmpeg_args = ffmpeg_args
self.oss_keys = oss_keys # 本次上传的临时音频 key(供事后清理)
def build_direct_render(
*,
resolved_clips: list[Any],
output_width: int,
output_height: int,
output_fps: int,
tts_audio: Optional[Path] = None,
bgm_audio: Optional[Path] = None,
title_text: str = "",
subtitle_segments: Optional[list[Any]] = None,
font: str = DEFAULT_DRAWTEXT_FONT,
vcodec: str = "h264_nvenc",
preset: str = "p4",
video_bitrate: str = "",
cq: int = 23,
edge_crop_pct: float = 0.0,
total_duration: float = 0.0,
) -> DirectRenderPlan:
"""构造 P4000 直连渲染所需的 inputs 与 ffmpeg_args。
视频:每段 trim/setpts/scale/pad/fps → concat(全硬切)→ 可选边缘 crop+scale → drawtext。
音频:concat 时丢弃原生音轨(只映射 [vfinal]),TTS/BGM 上传签名后 amix 混音。
"""
if not resolved_clips:
raise ValueError("build_direct_render: no resolved clips")
inputs: dict[str, str] = {}
oss_keys: list[str] = []
input_args: list[str] = []
fc: list[str] = [] # filter_complex 各段
n = len(resolved_clips)
# 1. 视频输入(原始素材签名 URL)
for i, clip in enumerate(resolved_clips):
sk = getattr(clip, "_storage_key", None)
if not sk:
raise ValueError(f"clip {getattr(clip, 'clip_id', i)} missing _storage_key")
fname = f"v{i}.mp4"
inputs[fname] = sign_asset_url(sk)
input_args.extend(["-i", fname])
# 2. 每个视频段预处理
pre_labels: list[str] = []
for i, clip in enumerate(resolved_clips):
filters: list[str] = []
start = float(getattr(clip, "start_time", 0) or 0)
eff = float(getattr(clip, "duration", 0) or 0)
if eff <= 0:
eff = float(getattr(clip, "actual_duration", 0) or 0)
if eff > 0:
if start > 0:
filters.append(f"trim=start={start:.3f}:duration={eff:.3f}")
else:
filters.append(f"trim=duration={eff:.3f}")
filters.append("setpts=PTS-STARTPTS")
speed = float(getattr(clip, "playback_speed", 1.0) or 1.0)
if abs(speed - 1.0) >= 1e-6:
filters.append(f"setpts=PTS/{speed:.4f}")
filters.append(f"scale={output_width}:{output_height}:force_original_aspect_ratio=decrease")
filters.append(f"pad={output_width}:{output_height}:trunc((ow-iw)/2):trunc((oh-ih)/2):black")
filters.append("setpts=PTS-STARTPTS")
filters.append(f"fps={output_fps}")
label = f"vc{i}"
fc.append(f"[{i}:v]{','.join(filters)}[{label}]")
pre_labels.append(label)
# 3. concat(全硬切;原生音频丢弃,v=1:a=0)
concat_in = "".join(f"[{l}]" for _lbl in pre_labels)
fc.append(f"{concat_in}concat=n={n}:v=1:a=0[vcat]")
cur = "vcat"
# 4. 边缘裁剪降重(合并进同一条,不再单独重编码)
if edge_crop_pct and edge_crop_pct > 0:
p = float(edge_crop_pct)
keep = 1.0 - 2.0 * p
cw_expr = f"trunc(iw*{keep:.4f}/2)*2"
ch_expr = f"trunc(ih*{keep:.4f}/2)*2"
fc.append(
f"[{cur}]crop=w='{cw_expr}':h='{ch_expr}':x='(iw-{cw_expr})/2':y='(ih-{ch_expr})/2',"
f"scale={output_width}:{output_height}[vcrop]"
)
cur = "vcrop"
# 5. drawtext 字幕(标题整段 + ASR 逐句)
draw_filters: list[str] = []
if title_text.strip():
title_size = max(int(output_height * 0.05), 24)
draw_filters.append(
build_drawtext_filter(
text=title_text,
start=0.0,
end=max(total_duration, 0.1),
font=font,
font_size=title_size,
y_expr="h-th-40",
box=True,
)
)
sub_size = max(int(output_height * 0.045), 20)
for seg in subtitle_segments or []:
txt = getattr(seg, "text", "") or ""
if not txt.strip():
continue
draw_filters.append(
build_drawtext_filter(
text=txt,
start=float(getattr(seg, "start", 0)),
end=float(getattr(seg, "end", 0)),
font=font,
font_size=sub_size,
y_expr="h-th-60",
borderw=2,
)
)
if draw_filters:
prev = cur
for idx, df in enumerate(draw_filters):
out_l = "vfinal" if idx == len(draw_filters) - 1 else f"vd{idx}"
fc.append(f"[{prev}]{df}[{out_l}]")
prev = out_l
vfinal_label = prev
else:
fc.append(f"[{cur}]format=yuv420p[vfinal]")
vfinal_label = "vfinal"
# 6. 音频输入与混音
audio_items: list[tuple[int, float]] = [] # (input_index, volume)
next_idx = n
if tts_audio and Path(tts_audio).exists():
turl, tkey = upload_local_audio_and_sign(Path(tts_audio))
tname = "tts" + (Path(tts_audio).suffix or ".mp3")
inputs[tname] = turl
oss_keys.append(tkey)
input_args.extend(["-i", tname])
audio_items.append((next_idx, 1.0))
next_idx += 1
if bgm_audio and Path(bgm_audio).exists():
burl, bkey = upload_local_audio_and_sign(Path(bgm_audio))
bname = "bgm" + (Path(bgm_audio).suffix or ".mp3")
inputs[bname] = burl
oss_keys.append(bkey)
input_args.extend(["-i", bname])
audio_items.append((next_idx, 0.35))
next_idx += 1
maps: list[str] = ["-map", f"[{vfinal_label}]"]
if audio_items:
mix_labels: list[str] = []
for k, (idx, vol) in enumerate(audio_items):
alabel = f"au{k}"
fc.append(
f"[{idx}:a]aresample=44100,volume={vol:.2f},aformat=sample_fmts=fltp:channel_layouts=stereo[{alabel}]"
)
mix_labels.append(alabel)
mix_in = "".join(f"[{l}]" for _lbl in mix_labels)
fc.append(f"{mix_in}amix=inputs={len(mix_labels)}:duration=first:dropout_transition=2,aresample=44100[afinal]")
maps.extend(["-map", "[afinal]", "-c:a", "aac", "-b:a", "128k"])
# 7. 组装 ffmpeg_args + NVENC 编码
ffmpeg_args = ["-y", *input_args, "-filter_complex", ";".join(fc), *maps]
ffmpeg_args.extend(["-c:v", vcodec, "-preset", preset, "-pix_fmt", "yuv420p"])
if video_bitrate:
ffmpeg_args.extend(["-b:v", video_bitrate])
else:
ffmpeg_args.extend(["-cq", str(cq)])
ffmpeg_args.extend(["-movflags", "+faststart", "-f", "mp4", "pipe:1"])
return DirectRenderPlan(inputs=inputs, ffmpeg_args=ffmpeg_args, oss_keys=oss_keys)
+13 -3
View File
@@ -189,7 +189,9 @@ class RenderAdapter:
self._report_progress(progress_cb, 15.0, f"下载素材({len(ready_clips)} 个)")
# 2. 下载素材
asset_path_map, rendered_clip_ids, failed_clip_ids = self._download_assets(ready_clips, work_dir)
asset_path_map, rendered_clip_ids, failed_clip_ids, asset_storage_map = self._download_assets(
ready_clips, work_dir
)
if not asset_path_map:
return RenderAdapterResult(
success=False,
@@ -206,6 +208,7 @@ class RenderAdapter:
plan=plan,
clips=ready_clips,
asset_path_map=asset_path_map,
asset_storage_map=asset_storage_map,
work_dir=work_dir,
plan_id=plan_id,
job_id=job_id,
@@ -315,7 +318,7 @@ class RenderAdapter:
def _download_assets(
self, clips: list[EditPlanClip], work_dir: Path
) -> tuple[dict[str, Path], list[str], list[str]]:
) -> tuple[dict[str, Path], list[str], list[str], dict[str, str]]:
"""下载片段素材到本地。
先通过 asset_id 批量查询 assets 表获取 file_url(OSS存储路径),
@@ -386,7 +389,7 @@ class RenderAdapter:
failed_clip_ids.append(clip.id)
logger.warning("素材下载失败: clip_id=%s asset_id=%s", clip.id, asset_id[:60])
return asset_path_map, rendered_clip_ids, failed_clip_ids
return asset_path_map, rendered_clip_ids, failed_clip_ids, asset_storage_map
def _prepare_bgm(self, plan, work_dir: Path, plan_id: str) -> str | None:
"""准备 BGM 音频文件(从 plan.config.bgm 读取配置)。
@@ -541,6 +544,7 @@ class RenderAdapter:
rendered_clip_ids: list[str] | None = None,
failed_clip_ids: list[str] | None = None,
voiceover_audio_path: str | None = None,
asset_storage_map: dict[str, str] | None = None,
) -> RenderAdapterResult:
"""执行统一渲染核心流程(BGM + ASR + 渲染 + 缩略图 + 上传)。
@@ -590,6 +594,12 @@ class RenderAdapter:
voiceover_audio_path=voiceover_audio_path,
clip_has_text=clip_has_text,
)
# 注入每个视频段对应素材的 storage_key,供全 GPU 直连管线直接签名下载
_storage_map = asset_storage_map or {}
for c in clips:
sk = _storage_map.get(getattr(c, "asset_id", ""))
if sk:
c._storage_key = sk
result = render_svc.render()
# 4.5 渲染后校验输出完整性
@@ -233,7 +233,7 @@ class UnifiedRenderService:
return
if abs(mt.brightness) > 1e-4 or abs(mt.contrast - 1.0) > 1e-4 or abs(mt.saturation - 1.0) > 1e-4:
filters.append(
f"eq=brightness={mt.brightness:+.4f}:" f"contrast={mt.contrast:.4f}:saturation={mt.saturation:.4f}"
f"eq=brightness={mt.brightness:+.4f}:contrast={mt.contrast:.4f}:saturation={mt.saturation:.4f}"
)
@staticmethod
@@ -341,6 +341,33 @@ class UnifiedRenderService:
len(pip_sources),
)
# 4.8 全 GPU 直连管线(P1):命中主流场景则跳过 mezzanine/边缘裁剪 CPU 重编码
output_path = self.work_dir / f"rendered_{self.plan.id}.mp4"
direct_result = self._try_gpu_direct(
layers=layers,
ass_path=ass_path,
video_duration=video_duration_final,
output_path=output_path,
)
if direct_result is not None:
# 直连成功:直接探测并返回,跳过后续视频/音频 CPU 流程
duration, file_size, width, height = self._probe_output(output_path)
logger.info(
"[unified-render] gpu-direct done: plan_id=%s total_ms=%d output_size=%d resolution=%dx%d",
self.plan.id,
int((time.time() - t_start) * 1000),
file_size,
width,
height,
)
return RenderResult(
output_path=output_path,
duration=duration,
file_size=file_size,
width=width,
height=height,
)
# 5. 视频主渲染
t_video_start = time.time()
video_only_path = self.work_dir / f"rendered_{self.plan.id}_video.mp4"
@@ -2180,6 +2207,170 @@ class UnifiedRenderService:
# ── GPU NVENC 加速 ────────────────────────────────────────────────────
# ── 全 GPU 直连渲染(P1)─────────────────────────────────────────────
def _can_use_gpu_direct(self, layers: list[RenderLayer]) -> bool:
"""判断是否命中直连支持的场景:单一主视频轨、全硬切、无复杂合成。"""
try:
cfg = self.plan.config or {}
# 特性开关(默认开启;可经 env/plan config 关闭灰度回退)
if not bool(cfg.get("gpu_direct_enabled", True)):
return False
video_layers = [_lyr for _lyr in layers if l.role not in ("audio",)]
# 只允许一个视频层,且角色为主层
if len(video_layers) != 1:
return False
role = video_layers[0].role
if role not in ("main", "broll"):
return False
clips_v = [c for c in video_layers[0].clips if c.clip_type != "audio"]
if not clips_v:
return False
# 全硬切(第一个 clip 的转场忽略)
for c in clips_v[1:]:
te = c.transition_effect
if te not in (None, "", "cut"):
return False
# 无画中画 / 水印 / 贴纸 / 片头片尾 / 绿幕 / 倒放 / 调色
if (cfg or {}).get("pip_config"):
return False
if (cfg or {}).get("intro_outro"):
return False
for c in clips_v:
cc = c.config or {}
if cc.get("watermark") or cc.get("stickers") or cc.get("chroma_key"):
return False
if ReverseConfig.from_dict(cc.get("reverse")).enabled:
return False
cg = ColorGradeConfig.from_dict(cc.get("color_grade"))
if cg.enabled and cg.has_effect():
return False
if not getattr(c, "_storage_key", None):
return False
return True
except Exception: # noqa: BLE001
logger.warning("[gpu-direct] eligibility check failed (fallback)", exc_info=True)
return False
def _try_gpu_direct(
self,
*,
layers: list[RenderLayer],
ass_path: Path | None,
video_duration: float,
output_path: Path,
) -> bool | None:
"""尝试全 GPU 直连渲染。成功返回 True,不支持/失败返回 None(调用方走旧链路)。"""
if not self._can_use_gpu_direct(layers):
return None
if not self._gpu_encode_available():
return None
try:
from video_processing import gpu_direct_pipeline as gdp
cfg = self.plan.config or {}
video_layer = next(_lyr for _lyr in layers if l.role not in ("audio",))
video_clips = [c for c in video_layer.clips if c.clip_type != "audio"]
# TTS:把 audio 层的 TTS 片段合并成一个文件给 P4000
tts_merged: Path | None = None
audio_layer = next((_lyr for _lyr in layers if l.role == "audio"), None)
if audio_layer:
tts_clips = [c for c in audio_layer.clips if (c.config or {}).get("tts") and c.local_path.exists()]
if tts_clips:
tts_merged = self._concat_audio_clips(tts_clips, tag="tts_direct")
# BGM 本地文件
bgm_path = Path(self.bgm_path) if self.bgm_path else None
if bgm_path is not None and not bgm_path.exists():
bgm_path = None
# 字幕:标题 + ASR 时间轴
title_cfg = cfg.get("title", {}) or cfg.get("title_config", {}) or {}
title_text = ""
if isinstance(title_cfg, dict) and title_cfg.get("enabled", True):
title_text = title_cfg.get("text", "") or ""
subtitle_segments: list[Any] = []
sub_cfg = cfg.get("subtitle", {}) or {}
if isinstance(sub_cfg, dict) and sub_cfg.get("enabled", True):
if sub_cfg.get("auto_generated") and self._asr_timeline_cache is not None:
subtitle_segments = list(self._asr_timeline_cache.segments)
# 边缘裁剪比例(与 random_edge_crop 默认 2~5% 同口径,取固定 3%)
dedup = self._dedup_enabled()
edge_pct = 0.03 if dedup else 0.0
plan = gdp.build_direct_render(
resolved_clips=video_clips,
output_width=self.output_width,
output_height=self.output_height,
output_fps=self.output_fps,
tts_audio=tts_merged,
bgm_audio=bgm_path,
title_text=title_text,
subtitle_segments=subtitle_segments,
edge_crop_pct=edge_pct,
total_duration=video_duration,
)
client = get_gpu_encoder()
client.render_inputs_to_output(plan.inputs, plan.ffmpeg_args, output_path)
# 清理本次上传的临时音频
for key in plan.oss_keys:
try:
from video_processing.oss_helpers import _storage
_storage().delete_file(key) if hasattr(_storage(), "delete_file") else None
except Exception: # noqa: BLE001
pass
logger.info("[gpu-direct] success: plan_id=%s clips=%d", self.plan.id, len(video_clips))
return True
except GpuEncodeError as e:
logger.warning("[gpu-direct] failed (fallback to legacy): %s", e)
try:
if output_path.exists():
output_path.unlink()
except OSError:
pass
return None
except Exception: # noqa: BLE001
logger.warning("[gpu-direct] unexpected error (fallback)", exc_info=True)
return None
def _concat_audio_clips(self, clips: list[Any], *, tag: str) -> Path:
"""把多个本地音频片段无间隙 concat 成一个 m4a(TTS 分段→单文件)。"""
out = self.work_dir / f"{tag}_{self.plan.id}.m4a"
listfile = self.work_dir / f"{tag}_{self.plan.id}.txt"
lines = []
for c in clips:
ap = str(c.local_path).replace("'", "'\\''")
lines.append(f"file '{ap}'")
listfile.write_text("\n".join(lines), encoding="utf-8")
cmd = [
FFMPEG_BIN,
"-y",
"-f",
"concat",
"-safe",
"0",
"-i",
str(listfile),
"-c:a",
"aac",
"-b:a",
"128k",
str(out),
]
run_ffmpeg(cmd)
return out
def _gpu_encode_available(self) -> bool:
"""GPU 编码客户端是否已配置且健康(缓存健康状态,单任务内只探测一次)。"""
if not getattr(self, "_gpu_health_ok", None):
@@ -2594,7 +2785,7 @@ class UnifiedRenderService:
b = pixel_pert.get("color_b", 0)
if r != 0 or g != 0 or b != 0:
# color_balance 参数范围 -1.0 ~ 1.0,这里用 /100 转换
filters.append(f"colorbalance=rs={r/100:.3f}:gs={g/100:.3f}:bs={b/100:.3f}")
filters.append(f"colorbalance=rs={r / 100:.3f}:gs={g / 100:.3f}:bs={b / 100:.3f}")
@staticmethod
def _clip_volume(clip: ResolvedClip) -> float:
+97
View File
@@ -309,6 +309,103 @@ class GpuEncoderClient:
except Exception as e: # noqa: BLE001
logger.warning("[gpu-encoder] failed to delete OSS mezzanine %s: %s", oss_key, e)
# ------------------------------------------------------------------
# High-level: render arbitrary inputs → final output (all-GPU pipeline)
# ------------------------------------------------------------------
def render_inputs_to_output(
self,
inputs: dict[str, str],
ffmpeg_args: list[str],
output_path: Path,
*,
timeout: Optional[int] = None,
) -> dict[str, Any]:
"""把多输入(原始素材/字幕/BGM)连同完整 filter_complex 交给 P4000 一次出片。
与 encode_mezzanine_to_output 的区别:worker 侧不再生成/上传 mezzanine,
P4000 直接从 inputs 中的签名 URL 下载原始素材,filter_complex 内完成
concat/scale/crop/drawtext/amix,末端 h264_nvenc 只编码一次。
传输:成片仍走 relay 回传(P4000 PUT → worker GET),避免公网 OSS 往返。
Args:
inputs: {裸文件名: 可下载URL},key 即 ffmpeg_args 中引用的文件名
ffmpeg_args: 完整 ffmpeg 参数(含 -i、-filter_complex、-map、NVENC 编码参数)
output_path: worker 本地成片落盘路径
timeout: P4000 侧超时(秒)
"""
if not inputs:
raise GpuEncodeError("render_inputs_to_output: inputs is empty")
if not ffmpeg_args:
raise GpuEncodeError("render_inputs_to_output: ffmpeg_args is empty")
if not self.relay_base_url:
raise GpuEncodeError("gpu_encode_relay_base_url not configured")
timeout = timeout or self.sync_timeout
t_total = time.time()
result_key: Optional[str] = None
try:
secret = self._get_relay_secret()
# 1. result relay URLs(成片 P4000 PUT → worker GET)
result_key = uuid.uuid4().hex
put_url = self._result_put_url(result_key, secret)
get_url = self._result_get_url(result_key, secret)
del_result_url = get_url
# 2. pre-warm then call P4000 sync render
self._warm_up_if_needed()
body = {
"inputs": dict(inputs),
"ffmpeg_args": list(ffmpeg_args),
"output_url": put_url,
"timeout": int(timeout),
}
job = self._post_sync(body)
self._last_ok_ts = time.time()
logger.info(
"[gpu-encoder] P4000 direct done: job_id=%s rc=%s size=%s dur=%ss inputs=%d",
job.get("job_id"),
job.get("ffmpeg_rc"),
job.get("size"),
job.get("duration"),
len(inputs),
)
# 3. download result from relay to output_path
output_path.parent.mkdir(parents=True, exist_ok=True)
size = self._download_to_file(get_url, output_path)
# 4. cleanup relay result
self._relay_delete(del_result_url)
logger.info(
"[gpu-encoder] direct render ok → %s (%d bytes) total=%.2fs",
output_path.name,
size,
time.time() - t_total,
)
return {
"job": job,
"output_size": size,
"output_path": str(output_path),
"transport": "direct",
}
except GpuEncodeError:
raise
except Exception as e: # noqa: BLE001
raise GpuEncodeError(f"unexpected: {e}") from e
finally:
# cleanup relay result (best-effort)
if result_key:
try:
secret = self._get_relay_secret()
self._relay_delete(self._relay_result_url(self.relay_internal_base_url, result_key, secret))
except Exception as e: # noqa: BLE001
logger.warning("[gpu-encoder] failed to delete relay result %s: %s", result_key, e)
# ------------------------------------------------------------------
# Internal helpers
# ------------------------------------------------------------------
+7 -7
View File
@@ -666,7 +666,7 @@ class TestDownloadAssets:
_make_clip("c2", order=1, asset_id="asset_002"),
]
asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path)
asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path)
assert len(asset_path_map) == 2
assert "asset_001" in asset_path_map
@@ -693,7 +693,7 @@ class TestDownloadAssets:
_make_clip("c2", order=1, asset_id="asset_002"),
]
asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path)
asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path)
assert len(asset_path_map) == 1
assert "asset_002" in asset_path_map
@@ -714,7 +714,7 @@ class TestDownloadAssets:
_make_clip("c1", order=0, asset_id="asset_001"),
]
asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path)
asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path)
assert len(asset_path_map) == 0
assert len(rendered_ids) == 0
@@ -750,7 +750,7 @@ class TestDownloadAssets:
_make_clip("c3", order=2, asset_id="asset_003"),
]
asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path)
asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path)
assert len(asset_path_map) == 2
assert "c1" in rendered_ids
@@ -771,7 +771,7 @@ class TestDownloadAssets:
_make_clip("c2", order=1, asset_id="asset_shared"),
]
asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path)
asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path)
assert len(asset_path_map) == 1
assert mock_download.call_count == 1
@@ -795,7 +795,7 @@ class TestDownloadAssets:
adapter = RenderAdapter(mock_db)
clips = [_make_clip("c1", order=0, asset_id="asset_fallback")]
asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path)
asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path)
assert len(asset_path_map) == 1
assert "c1" in rendered_ids
@@ -820,7 +820,7 @@ class TestDownloadAssets:
adapter = RenderAdapter(mock_db)
clips = [_make_clip("c1", order=0, asset_id="asset_no_key")]
asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path)
asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path)
assert len(asset_path_map) == 0
assert "c1" in failed_ids