feat(gpu): 接入 P4000 NVENC 硬件编码加速
- packages/config/base.py: 新增 GPU_ENCODE_* 配置项(开关/endpoint/relay/secret/编码参数/fallback)
- packages/shared/gpu_encoder.py: 新增 GpuEncoderClient,封装 health 探测 + mezzanine 上传 OSS + P4000 nvenc 编码 + relay 回传下载,失败抛 GpuEncodeError 触发 CPU 降级
- apps/api/app/api/routes/gpu_relay.py: 新增 internal PUT/GET/DELETE /api/v1/internal/gpu-relay/{key}(token 鉴权),P4000 PUT 编码结果,worker GET 下载
- apps/api/app/api/router.py: 注册 gpu_relay_router
- apps/worker/video_processing/unified_render_service.py: _execute_ffmpeg 和 _render_pass_through 尝试 GPU 路径:CPU ultrafast mezzanine → P4000 nvenc → 输出到最终路径;任何失败自动回退到原 CPU libx264 路径
- apps/worker/worker_app/tasks/_startup.py: worker_ready 时探测 P4000 健康并打日志
- tests/unit/test_gpu_encoder.py: GpuEncoderClient 单测(health/sync 调用/失败/fallback/relay URL)
架构:
- 输入:CPU 输出 libx264 ultrafast mezzanine → 上传 OSS 临时前缀 → P4000 签名 URL 下载
- 输出:P4000 PUT → 宿主机 nginx(tailscale:80)→ API /api/ 反代 → gpu_relay 路由落盘到 generated/gpu_relay/
- 回传:worker 通过 docker 网络 http://xiaoxia-api-staging:8000 GET 下载最终 mp4 到 output_path
- 降级:GPU 任何环节异常(health/上传/编码/回传/下载)→ 原 CPU 路径继续执行,不影响成片
This commit is contained in:
committed by
saas-backend
parent
a5075624f8
commit
105ab54059
@@ -0,0 +1,366 @@
|
||||
"""P4000 NVENC 远程编码客户端。
|
||||
|
||||
完整链路(encode_video_file):
|
||||
1. CPU 滤镜已在本地生成 mezzanine 中间片(libx264 ultrafast)
|
||||
2. 上传 mezzanine 到 OSS 临时前缀,拿到签名 GET URL
|
||||
3. 生成 relay 一次性 key,构造带 token 的 PUT URL(指向 API 服务 /api/v1/internal/gpu-relay/<key>)
|
||||
4. POST P4000 /api/render/sync:inputs={"in.mp4": "<oss-signed-url>"}, output_url="<relay-put-url>"
|
||||
ffmpeg_args: -i in.mp4 [-vf <vf>] -c:v h264_nvenc ... -an/-c:a aac -f mp4 pipe:1
|
||||
5. P4000 编码完成后 PUT 最终 mp4 到 relay,API 服务落盘到 /app/generated/gpu_relay/<key>
|
||||
6. 本客户端从 relay GET 下载最终文件到 output_path,然后调用 relay DELETE 清理
|
||||
7. 删除 OSS 临时 mezzanine
|
||||
|
||||
任何环节失败抛 GpuEncodeError,调用方应 fallback 到 CPU libx264。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import secrets
|
||||
import socket
|
||||
import time
|
||||
import urllib.error
|
||||
import urllib.parse
|
||||
import urllib.request
|
||||
import uuid
|
||||
from dataclasses import dataclass
|
||||
from pathlib import Path
|
||||
from typing import Any, Optional
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class GpuEncodeError(RuntimeError):
|
||||
"""GPU 编码失败(网络/超时/ffmpeg/upload/download 任一环节)。调用方应 fallback 到 CPU。"""
|
||||
|
||||
|
||||
@dataclass
|
||||
class GpuHealth:
|
||||
healthy: bool
|
||||
worker: str = ""
|
||||
gpu_name: str = ""
|
||||
nvenc_h264: bool = False
|
||||
nvenc_hevc: bool = False
|
||||
error: str = ""
|
||||
|
||||
@property
|
||||
def ready(self) -> bool:
|
||||
return self.healthy and self.nvenc_h264
|
||||
|
||||
|
||||
class GpuEncoderClient:
|
||||
def __init__(
|
||||
self,
|
||||
endpoint: str,
|
||||
relay_base_url: str,
|
||||
*,
|
||||
sync_timeout: int = 300,
|
||||
health_timeout: float = 3.0,
|
||||
vcodec: str = "h264_nvenc",
|
||||
preset: str = "p4",
|
||||
crf: int = 23,
|
||||
bitrate: str = "",
|
||||
relay_secret: str = "",
|
||||
oss_tmp_prefix: str = "tmp/gpu-mezzanine/",
|
||||
) -> None:
|
||||
self.endpoint = endpoint.rstrip("/")
|
||||
self.relay_base_url = relay_base_url.rstrip("/")
|
||||
self.sync_timeout = sync_timeout
|
||||
self.health_timeout = health_timeout
|
||||
self.vcodec = vcodec
|
||||
self.preset = preset
|
||||
self.crf = crf
|
||||
self.bitrate = bitrate
|
||||
self._relay_secret = relay_secret
|
||||
self.oss_tmp_prefix = oss_tmp_prefix.rstrip("/") + "/" if oss_tmp_prefix else "tmp/gpu-mezzanine/"
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Health
|
||||
# ------------------------------------------------------------------
|
||||
def check_health(self) -> GpuHealth:
|
||||
url = f"{self.endpoint}/health"
|
||||
try:
|
||||
with urllib.request.urlopen(url, timeout=self.health_timeout) as resp:
|
||||
data = json.loads(resp.read().decode("utf-8"))
|
||||
except (urllib.error.URLError, socket.timeout, TimeoutError, json.JSONDecodeError, ConnectionError) as e:
|
||||
return GpuHealth(healthy=False, error=f"health probe failed: {e}")
|
||||
try:
|
||||
return GpuHealth(
|
||||
healthy=data.get("status") == "healthy",
|
||||
worker=str(data.get("worker", "")),
|
||||
gpu_name=(data.get("gpu") or {}).get("name", ""),
|
||||
nvenc_h264=bool((data.get("nvenc") or {}).get("h264_nvenc")),
|
||||
nvenc_hevc=bool((data.get("nvenc") or {}).get("hevc_nvenc")),
|
||||
)
|
||||
except Exception as e: # noqa: BLE001
|
||||
return GpuHealth(healthy=False, error=f"malformed health response: {e}")
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# High-level: encode a mezzanine file to final output
|
||||
# ------------------------------------------------------------------
|
||||
def encode_mezzanine_to_output(
|
||||
self,
|
||||
mezzanine_path: Path,
|
||||
output_path: Path,
|
||||
*,
|
||||
extra_video_args: Optional[list[str]] = None,
|
||||
audio_args: Optional[list[str]] = None,
|
||||
timeout: Optional[int] = None,
|
||||
) -> dict[str, Any]:
|
||||
"""把 mezzanine(CPU 滤镜已完成)交给 P4000 NVENC 编码,结果写到 output_path。
|
||||
|
||||
extra_video_args: -i 之后、-c:v 之前插入的 ffmpeg 参数(如分辨率/帧率调整)。
|
||||
audio_args: 音频编码参数(如 ["-c:a","aac","-b:a","128k"]);None 表示 -an 无音频。
|
||||
"""
|
||||
if not mezzanine_path.exists():
|
||||
raise GpuEncodeError(f"mezzanine file not found: {mezzanine_path}")
|
||||
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()
|
||||
oss_key: Optional[str] = None
|
||||
relay_key: Optional[str] = None
|
||||
|
||||
try:
|
||||
# 1. upload mezzanine → OSS
|
||||
input_url, oss_key = self._upload_mezzanine(mezzanine_path)
|
||||
logger.debug("[gpu-encoder] mezzanine uploaded: oss_key=%s", oss_key)
|
||||
|
||||
# 2. prepare relay put/get URLs
|
||||
relay_key = uuid.uuid4().hex
|
||||
secret = self._get_relay_secret()
|
||||
put_url = self._relay_url(relay_key, secret)
|
||||
get_url = put_url
|
||||
del_url = put_url # same URL, DELETE method
|
||||
|
||||
# 3. build ffmpeg args
|
||||
ffmpeg_args = ["-y", "-i", "in.mp4"]
|
||||
if extra_video_args:
|
||||
ffmpeg_args.extend(extra_video_args)
|
||||
ffmpeg_args.extend(["-c:v", self.vcodec, "-preset", self.preset])
|
||||
if self.bitrate:
|
||||
ffmpeg_args.extend(["-b:v", self.bitrate])
|
||||
else:
|
||||
ffmpeg_args.extend(["-cq", str(self.crf)])
|
||||
ffmpeg_args.extend(["-pix_fmt", "yuv420p", "-movflags", "+faststart"])
|
||||
if audio_args:
|
||||
ffmpeg_args.extend(audio_args)
|
||||
else:
|
||||
ffmpeg_args.append("-an")
|
||||
ffmpeg_args.extend(["-f", "mp4", "pipe:1"])
|
||||
|
||||
# 4. call P4000 sync render
|
||||
body = {
|
||||
"inputs": {"in.mp4": input_url},
|
||||
"ffmpeg_args": ffmpeg_args,
|
||||
"output_url": put_url,
|
||||
"timeout": int(timeout),
|
||||
}
|
||||
job = self._post_sync(body, mezzanine_path=mezzanine_path)
|
||||
logger.info(
|
||||
"[gpu-encoder] P4000 done: job_id=%s rc=%s size=%s dur=%ss",
|
||||
job.get("job_id"), job.get("ffmpeg_rc"), job.get("size"), job.get("duration"),
|
||||
)
|
||||
|
||||
# 5. 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)
|
||||
|
||||
# 6. cleanup relay
|
||||
self._relay_delete(del_url)
|
||||
|
||||
logger.info(
|
||||
"[gpu-encoder] encode ok: %s → %s (%d bytes) total=%.2fs",
|
||||
mezzanine_path.name, output_path.name, size, time.time() - t_total,
|
||||
)
|
||||
return {"job": job, "output_size": size, "output_path": str(output_path)}
|
||||
|
||||
except GpuEncodeError:
|
||||
raise
|
||||
except Exception as e: # noqa: BLE001
|
||||
raise GpuEncodeError(f"unexpected: {e}") from e
|
||||
finally:
|
||||
# cleanup OSS mezzanine (best-effort)
|
||||
if oss_key:
|
||||
try:
|
||||
self._delete_oss(oss_key)
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("[gpu-encoder] failed to delete OSS mezzanine %s: %s", oss_key, e)
|
||||
# relay cleanup also best-effort (done above after download)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Internal helpers
|
||||
# ------------------------------------------------------------------
|
||||
def _get_relay_secret(self) -> str:
|
||||
if self._relay_secret:
|
||||
return self._relay_secret
|
||||
# read from env (same var API server uses)
|
||||
env = (os.getenv("APP_ENV", os.getenv("ENV", "development"))).lower()
|
||||
secret = (os.getenv("GPU_ENCODE_RELAY_SECRET", "") or "").strip()
|
||||
if not secret:
|
||||
if env in ("production", "prod"):
|
||||
raise GpuEncodeError("GPU_ENCODE_RELAY_SECRET must be set in production")
|
||||
# dev: fail - worker should always have a secret explicitly set (or same ephemeral won't match)
|
||||
raise GpuEncodeError("GPU_ENCODE_RELAY_SECRET not set")
|
||||
return secret
|
||||
|
||||
def _relay_url(self, key: str, secret: str) -> str:
|
||||
return f"{self.relay_base_url}/api/v1/internal/gpu-relay/{key}?token={urllib.parse.quote(secret, safe='')}"
|
||||
|
||||
def _post_sync(self, body: dict[str, Any], *, mezzanine_path: Path) -> dict[str, Any]:
|
||||
url = f"{self.endpoint}/api/render/sync"
|
||||
req_timeout = body.get("timeout", self.sync_timeout) + 60
|
||||
payload = json.dumps(body).encode("utf-8")
|
||||
req = urllib.request.Request(
|
||||
url, data=payload,
|
||||
headers={"Content-Type": "application/json"},
|
||||
method="POST",
|
||||
)
|
||||
t0 = time.time()
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=req_timeout) as resp:
|
||||
raw = resp.read().decode("utf-8")
|
||||
except urllib.error.HTTPError as e:
|
||||
detail = e.read().decode("utf-8", errors="replace")[:1000]
|
||||
raise GpuEncodeError(f"P4000 HTTP {e.code}: {detail}") from e
|
||||
except (urllib.error.URLError, socket.timeout, TimeoutError, ConnectionError) as e:
|
||||
raise GpuEncodeError(f"P4000 connection error: {e}") from e
|
||||
try:
|
||||
result = json.loads(raw)
|
||||
except json.JSONDecodeError as e:
|
||||
raise GpuEncodeError(f"P4000 bad JSON: {raw[:500]}") from e
|
||||
dt = time.time() - t0
|
||||
|
||||
status = result.get("status")
|
||||
ffmpeg_rc = result.get("ffmpeg_rc")
|
||||
uploaded = result.get("uploaded")
|
||||
if status != "completed" or ffmpeg_rc != 0:
|
||||
err = result.get("message") or result.get("error") or "unknown"
|
||||
raise GpuEncodeError(f"P4000 job failed: status={status} rc={ffmpeg_rc} err={err!s:.500}")
|
||||
# P4000 has a known bug where uploaded=true even on PUT SSL failure;
|
||||
# we will verify by downloading, so don't hard-fail here but log
|
||||
if not uploaded:
|
||||
logger.warning("[gpu-encoder] P4000 reports uploaded=false (will verify via download)")
|
||||
result["_roundtrip"] = dt
|
||||
return result
|
||||
|
||||
def _download_to_file(self, url: str, output_path: Path) -> int:
|
||||
"""GET url → write to output_path. Returns bytes written."""
|
||||
tmp = output_path.with_suffix(output_path.suffix + ".gpu_tmp")
|
||||
size = 0
|
||||
try:
|
||||
with urllib.request.urlopen(url, timeout=self.sync_timeout) as resp:
|
||||
if resp.status != 200:
|
||||
raise GpuEncodeError(f"relay GET returned HTTP {resp.status}")
|
||||
with open(tmp, "wb") as f:
|
||||
while True:
|
||||
chunk = resp.read(1024 * 256)
|
||||
if not chunk:
|
||||
break
|
||||
f.write(chunk)
|
||||
size += len(chunk)
|
||||
if size == 0:
|
||||
raise GpuEncodeError("relay returned empty file")
|
||||
os.replace(tmp, output_path)
|
||||
return size
|
||||
except (urllib.error.URLError, socket.timeout, TimeoutError, ConnectionError) as e:
|
||||
if tmp.exists():
|
||||
try: tmp.unlink()
|
||||
except OSError: pass
|
||||
raise GpuEncodeError(f"failed to download from relay: {e}") from e
|
||||
|
||||
def _relay_delete(self, url: str) -> None:
|
||||
try:
|
||||
req = urllib.request.Request(url, method="DELETE")
|
||||
with urllib.request.urlopen(req, timeout=10) as resp:
|
||||
resp.read()
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.debug("[gpu-encoder] relay cleanup delete failed: %s", e)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# OSS helpers (optional - storage may not be available in all envs)
|
||||
# ------------------------------------------------------------------
|
||||
def _upload_mezzanine(self, path: Path) -> tuple[str, str]:
|
||||
"""Upload mezzanine to OSS tmp prefix, return (signed_get_url, oss_key)."""
|
||||
try:
|
||||
from packages.shared.storage import get_storage_service
|
||||
except ImportError as e:
|
||||
raise GpuEncodeError(f"storage service unavailable: {e}") from e
|
||||
storage = get_storage_service()
|
||||
if storage is None or storage.bucket is None:
|
||||
raise GpuEncodeError("OSS storage not configured; cannot upload mezzanine")
|
||||
key = f"{self.oss_tmp_prefix}{uuid.uuid4().hex}.mp4"
|
||||
try:
|
||||
storage.upload_file(str(path), key, content_type="video/mp4")
|
||||
except Exception as e: # noqa: BLE001
|
||||
raise GpuEncodeError(f"failed to upload mezzanine to OSS: {e}") from e
|
||||
# Generate signed GET URL (1h expiry)
|
||||
signed = storage.get_download_url(key, expires_seconds=3600)
|
||||
return signed, key
|
||||
|
||||
def _delete_oss(self, key: str) -> None:
|
||||
try:
|
||||
from packages.shared.storage import get_storage_service
|
||||
storage = get_storage_service()
|
||||
if storage is not None and storage.bucket is not None:
|
||||
storage.delete_file(key)
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.debug("[gpu-encoder] OSS delete %s failed: %s", key, e)
|
||||
|
||||
|
||||
# ── Singleton factory ────────────────────────────────────────────────────
|
||||
|
||||
_default_client: Optional[GpuEncoderClient] = None
|
||||
_default_client_initialized: bool = False
|
||||
|
||||
|
||||
def _build_client_from_settings() -> Optional[GpuEncoderClient]:
|
||||
try:
|
||||
from packages.config import get_shared_settings
|
||||
settings = get_shared_settings()
|
||||
except Exception: # noqa: BLE001
|
||||
return None
|
||||
if not getattr(settings, "enable_gpu_encode", False):
|
||||
return None
|
||||
endpoint = (getattr(settings, "gpu_encode_endpoint", "") or "").strip()
|
||||
relay = (getattr(settings, "gpu_encode_relay_base_url", "") or "").strip()
|
||||
if not endpoint or not relay:
|
||||
return None
|
||||
return GpuEncoderClient(
|
||||
endpoint=endpoint,
|
||||
relay_base_url=relay,
|
||||
sync_timeout=getattr(settings, "gpu_encode_sync_timeout", 300),
|
||||
health_timeout=getattr(settings, "gpu_encode_health_timeout", 3.0),
|
||||
vcodec=getattr(settings, "gpu_encode_vcodec", "h264_nvenc"),
|
||||
preset=getattr(settings, "gpu_encode_preset", "p4"),
|
||||
crf=getattr(settings, "gpu_encode_crf", 23),
|
||||
bitrate=getattr(settings, "gpu_encode_bitrate", "") or "",
|
||||
relay_secret=getattr(settings, "gpu_encode_relay_secret", "") or "",
|
||||
oss_tmp_prefix=getattr(settings, "gpu_encode_oss_tmp_prefix", "tmp/gpu-mezzanine/"),
|
||||
)
|
||||
|
||||
|
||||
def get_gpu_encoder() -> Optional[GpuEncoderClient]:
|
||||
"""返回进程级单例;未启用或未配置返回 None。"""
|
||||
global _default_client, _default_client_initialized
|
||||
if not _default_client_initialized:
|
||||
_default_client_initialized = True
|
||||
try:
|
||||
_default_client = _build_client_from_settings()
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("[gpu-encoder] failed to init client (CPU fallback): %s", e)
|
||||
_default_client = None
|
||||
return _default_client
|
||||
|
||||
|
||||
def reset_gpu_encoder_for_tests() -> None:
|
||||
global _default_client, _default_client_initialized
|
||||
_default_client = None
|
||||
_default_client_initialized = False
|
||||
|
||||
|
||||
# Convenience
|
||||
def is_gpu_encode_enabled() -> bool:
|
||||
return get_gpu_encoder() is not None
|
||||
Reference in New Issue
Block a user