perf(gpu): mezzanine relay直传跳过公网OSS(-18s), preset veryfast(-38%), 修复CI staging IP (#2063)
CI/CD Pipeline / Check if frontend-only change (push) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (push) Successful in 1s
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 2s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
CI/CD Pipeline / Frontend Lint (push) Has been skipped
CI/CD Pipeline / PR Build API Image (push) Has been skipped
CI/CD Pipeline / PR Build Web Image (push) Has been skipped
CI/CD Pipeline / PR Build Worker Image (push) Has been skipped
CI/CD Pipeline / Validate - Style (pull_request) Has been skipped
CI/CD Pipeline / Validate - Security (pull_request) Has been skipped
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Has been skipped
CI/CD Pipeline / Build Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Check push changed paths (push) Successful in 13s
CI/CD Pipeline / Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 55s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 1m14s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (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 / Build Staging API Image (push) Successful in 1m10s
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 / ACR Image Cleanup (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 40s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m31s
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 2m56s
CI/CD Pipeline / Build Staging Web Image (push) Successful in 2m24s
CI/CD Pipeline / Retag skipped Staging API Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Has been skipped
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m27s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been skipped
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Failing after 33s
CI/CD Pipeline / Staging E2E Tests (push) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (push) Has been skipped
CI/CD Pipeline / ACR Image Cleanup (push) Has been skipped
CI/CD Pipeline / Validate - Style (push) Successful in 3m53s
CI/CD Pipeline / Frontend Unit Tests (push) Failing after 4m47s
CI/CD Pipeline / Integration Tests (push) Successful in 5m28s
AI Code Review / AI Code Review (pull_request) Successful in 6m51s
CI/CD Pipeline / Validate - Security (push) Successful in 7m4s
CI/CD Pipeline / PR Build Web Image (pull_request) Failing after 9m19s
CI/CD Pipeline / CI Gate (pull_request) Failing after 1s
CI/CD Pipeline / Unit Tests (push) Successful in 10m32s
CI/CD Pipeline / Build Production API Image (push) Has been skipped
CI/CD Pipeline / CI Gate (push) Has been skipped
CI/CD Pipeline / Build Production Worker Image (push) Has been skipped
CI/CD Pipeline / Build Production Web Image (push) Has been skipped
CI/CD Pipeline / Canary Release to Production (push) Has been skipped
CI/CD Pipeline / Deploy Production (push) Has been skipped
CI/CD Pipeline / Production Browser E2E (push) Has been skipped

Co-authored-by: xiaoxia <dev@xiaoxiajianji.com>
Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
This commit was merged in pull request #2063.
This commit is contained in:
2026-09-27 14:48:54 +08:00
committed by auto-approve-bot
parent 9c0535da19
commit e4723bfb1b
9 changed files with 429 additions and 177 deletions
+124 -41
View File
@@ -2,15 +2,16 @@
完整链路(encode_video_file):
1. CPU 滤镜已在本地生成 mezzanine 中间片(libx264 ultrafast)
2. 上传 mezzanine 到 OSS 临时前缀,拿到签名 GET URL
2. 通过 HTTP PUT 把 mezzanine 上传到 relay(走 Tailscale/Docker 内网,~1s 完成)
- 失败则 fallback 到 OSS 上传(旧路径,兼容没有 :8092 内网可达的环境)
3. 生成 relay 一次性 key,构造两个带 token 的 URL:
- put_url:给 P4000 回传结果,走 relay_base_url(外部可达,通常是 host:port 经 nginx)
- put_url:给 P4000 回传结果,走 relay_base_url(Tailscale host:8092)
- get/del_url:worker 自己下载+清理用,走 relay_internal_base_url(Docker DNS 直连 API)
4. POST P4000 /api/render/sync:inputs={"in.mp4": "<oss-signed-url>"}, output_url="<put_url>"
4. POST P4000 /api/render/sync:inputs={"in.mp4": "<mezzanine-get-url>"}, output_url="<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 到 put_url,API 服务落盘到 /app/generated/gpu_relay/<key>
5. P4000 从 relay GET mezzanine → h264_nvenc 编码 → PUT 最终 mp4 到 put_url
6. 本客户端通过 get_url(Docker 内网)下载最终文件到 output_path,然后 DELETE 清理
7. 删除 OSS 临时 mezzanine
7. 删除 relay 上的 mezzanine 临时文件(以及 OSS fallback 的 key)
任何环节失败抛 GpuEncodeError,调用方应 fallback 到 CPU libx264。
"""
@@ -58,6 +59,8 @@ class GpuEncoderClient:
relay_base_url: str,
*,
relay_internal_base_url: str = "",
# Mezzanine 上传:默认走 relay(Tailscale/Docker 内网);设为 "oss" 强制走旧 OSS 路径
mezzanine_transport: str = "relay",
sync_timeout: int = 300,
health_timeout: float = 3.0,
vcodec: str = "h264_nvenc",
@@ -74,6 +77,7 @@ class GpuEncoderClient:
self.relay_internal_base_url = (
relay_internal_base_url.rstrip("/") if relay_internal_base_url else self.relay_base_url
)
self.mezzanine_transport = mezzanine_transport.lower() # "relay" | "oss"
self.sync_timeout = sync_timeout
self.health_timeout = health_timeout
self.vcodec = vcodec
@@ -88,16 +92,30 @@ class GpuEncoderClient:
# ------------------------------------------------------------------
# URL builders
# ------------------------------------------------------------------
def _relay_url_from_base(self, base_url: str, key: str, secret: str) -> str:
return f"{base_url}{self.RELAY_PATH_PREFIX}/{key}?token={urllib.parse.quote(secret, safe='')}"
def _relay_url_from_base(self, base_url: str, path: str, key: str, secret: str) -> str:
return f"{base_url}{self.RELAY_PATH_PREFIX}{path}/{key}?token={urllib.parse.quote(secret, safe='')}"
def _relay_put_url(self, key: str, secret: str) -> str:
"""给 P4000 回传结果用的 URL(外部可达)。"""
return self._relay_url_from_base(self.relay_base_url, key, secret)
def _relay_result_url(self, base_url: str, key: str, secret: str) -> str:
return self._relay_url_from_base(base_url, "", key, secret)
def _relay_internal_url(self, key: str, secret: str) -> str:
"""Worker 自己 GET/DELETE 用的 URL(Docker 内网)。"""
return self._relay_url_from_base(self.relay_internal_base_url, key, secret)
def _relay_mezz_url(self, base_url: str, key: str, secret: str) -> str:
return self._relay_url_from_base(base_url, "/mezzanine", key, secret)
def _result_put_url(self, key: str, secret: str) -> str:
"""P4000 回传编码结果 PUT URL(外部/Tailscale 可达)。"""
return self._relay_result_url(self.relay_base_url, key, secret)
def _result_get_url(self, key: str, secret: str) -> str:
"""Worker 下载最终结果 GET URL(Docker 内网)。"""
return self._relay_result_url(self.relay_internal_base_url, key, secret)
def _mezz_put_url(self, key: str, secret: str) -> str:
"""Worker 上传 mezzanine PUT URL(Docker 内网,快)。"""
return self._relay_mezz_url(self.relay_internal_base_url, key, secret)
def _mezz_get_url_for_p4000(self, key: str, secret: str) -> str:
"""P4000 下载 mezzanine GET URL(必须是 P4000 可达地址,Tailscale host:8092)。"""
return self._relay_mezz_url(self.relay_base_url, key, secret)
# ------------------------------------------------------------------
# Health
@@ -134,8 +152,8 @@ class GpuEncoderClient:
) -> 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 无音频。
传输:默认通过 relay PUT/GET(Tailscale 内网,省去公网 OSS 往返 18-20s);
若 relay PUT 失败且配置可用,自动 fallback 到 OSS。
"""
if not mezzanine_path.exists():
raise GpuEncodeError(f"mezzanine file not found: {mezzanine_path}")
@@ -145,19 +163,53 @@ class GpuEncoderClient:
timeout = timeout or self.sync_timeout
t_total = time.time()
oss_key: Optional[str] = None
relay_key: Optional[str] = None
mezz_key: Optional[str] = None
result_key: Optional[str] = None
input_url: str = ""
used_transport = self.mezzanine_transport
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 URLs (PUT 走外部 URL 给 P4000;GET/DELETE 走内部 Docker 网络)
relay_key = uuid.uuid4().hex
secret = self._get_relay_secret()
put_url = self._relay_put_url(relay_key, secret)
get_url = self._relay_internal_url(relay_key, secret)
del_url = get_url # 内部 URL,DELETE method
# 1. 上传 mezzanine 到 relay(或 OSS fallback)
mezz_size = mezzanine_path.stat().st_size
if used_transport == "relay":
mezz_key = uuid.uuid4().hex
mezz_put = self._mezz_put_url(mezz_key, secret)
mezz_get_for_p4000 = self._mezz_get_url_for_p4000(mezz_key, secret)
t_up = time.time()
try:
self._upload_file_put(mezz_put, mezzanine_path, "video/mp4")
input_url = mezz_get_for_p4000
logger.info(
"[gpu-encoder] mezzanine uploaded to relay: key=%s size=%d took=%.2fs",
mezz_key,
mezz_size,
time.time() - t_up,
)
except (GpuEncodeError, OSError, urllib.error.URLError) as e:
logger.warning(
"[gpu-encoder] relay mezz upload failed (%s), fallback to OSS",
e,
)
used_transport = "oss"
# relay 上传失败的部分文件 best-effort 清理
if mezz_key:
try:
self._relay_delete(self._relay_mezz_url(self.relay_internal_base_url, mezz_key, secret))
except Exception: # noqa: BLE001
pass
mezz_key = None
if used_transport == "oss" or not input_url:
input_url, oss_key = self._upload_mezzanine_to_oss(mezzanine_path)
logger.debug("[gpu-encoder] mezzanine uploaded to OSS: oss_key=%s", oss_key)
# 2. prepare result relay URLs
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
# 3. build ffmpeg args
ffmpeg_args = ["-y", "-i", "in.mp4"]
@@ -182,43 +234,56 @@ class GpuEncoderClient:
"output_url": put_url,
"timeout": int(timeout),
}
job = self._post_sync(body, mezzanine_path=mezzanine_path)
job = self._post_sync(body)
logger.info(
"[gpu-encoder] P4000 done: job_id=%s rc=%s size=%s dur=%ss",
"[gpu-encoder] P4000 done: job_id=%s rc=%s size=%s dur=%ss transport=%s",
job.get("job_id"),
job.get("ffmpeg_rc"),
job.get("size"),
job.get("duration"),
used_transport,
)
# 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)
# 6. cleanup relay result
self._relay_delete(del_result_url)
logger.info(
"[gpu-encoder] encode ok: %s → %s (%d bytes) total=%.2fs",
"[gpu-encoder] encode ok: %s → %s (%d bytes) total=%.2fs transport=%s",
mezzanine_path.name,
output_path.name,
size,
time.time() - t_total,
used_transport,
)
return {"job": job, "output_size": size, "output_path": str(output_path)}
return {
"job": job,
"output_size": size,
"output_path": str(output_path),
"transport": used_transport,
}
except GpuEncodeError:
raise
except Exception as e: # noqa: BLE001
raise GpuEncodeError(f"unexpected: {e}") from e
finally:
# cleanup relay mezzanine (best-effort)
if mezz_key:
try:
secret = self._get_relay_secret()
self._relay_delete(self._relay_mezz_url(self.relay_internal_base_url, mezz_key, secret))
except Exception as e: # noqa: BLE001
logger.warning("[gpu-encoder] failed to delete relay mezzanine %s: %s", mezz_key, e)
# 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
@@ -226,17 +291,15 @@ class GpuEncoderClient:
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 _post_sync(self, body: dict[str, Any], *, mezzanine_path: Path) -> dict[str, Any]:
def _post_sync(self, body: dict[str, Any]) -> 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")
@@ -267,8 +330,6 @@ class GpuEncoderClient:
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
@@ -309,10 +370,33 @@ class GpuEncoderClient:
except Exception as e: # noqa: BLE001
logger.debug("[gpu-encoder] relay cleanup delete failed: %s", e)
def _upload_file_put(self, url: str, path: Path, content_type: str) -> None:
"""HTTP PUT 流式上传文件到指定 URL(mezzanine 上传到 relay 用,Tailscale/Docker 内网)。"""
file_size = path.stat().st_size
# 使用生成器/文件对象流式上传,避免一次性载入大 mezzanine 文件到内存
with open(path, "rb") as f:
req = urllib.request.Request(
url,
data=f, # 文件对象支持read(),urllib会流式发送(但需要Content-Length)
method="PUT",
headers={
"Content-Type": content_type,
"Content-Length": str(file_size),
},
)
# 超时:按 ~20MB/s 内网速度估算 + 30s 保底
put_timeout = max(60, int(file_size / (20 * 1024 * 1024)) + 30)
with urllib.request.urlopen(req, timeout=put_timeout) as resp:
if resp.status not in (200, 201, 204):
body = resp.read().decode("utf-8", errors="replace")[:500]
raise GpuEncodeError(f"relay PUT failed: HTTP {resp.status} {body}")
resp.read()
# ------------------------------------------------------------------
# OSS helpers (optional - storage may not be available in all envs)
# OSS helpers (fallback)
# ------------------------------------------------------------------
def _upload_mezzanine(self, path: Path) -> tuple[str, str]:
def _upload_mezzanine_to_oss(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
@@ -326,7 +410,6 @@ class GpuEncoderClient:
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
@@ -365,6 +448,7 @@ def _build_client_from_settings() -> Optional[GpuEncoderClient]:
endpoint=endpoint,
relay_base_url=relay,
relay_internal_base_url=relay_internal,
mezzanine_transport=getattr(settings, "gpu_encode_mezzanine_transport", "relay") or "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"),
@@ -395,6 +479,5 @@ def reset_gpu_encoder_for_tests() -> None:
_default_client_initialized = False
# Convenience
def is_gpu_encode_enabled() -> bool:
return get_gpu_encoder() is not None