From e4723bfb1bd613bd6600a4560497273259a3072b Mon Sep 17 00:00:00 2001 From: xiaoxia Date: Sun, 27 Sep 2026 14:48:54 +0800 Subject: [PATCH] =?UTF-8?q?perf(gpu):=20mezzanine=20relay=E7=9B=B4?= =?UTF-8?q?=E4=BC=A0=E8=B7=B3=E8=BF=87=E5=85=AC=E7=BD=91OSS(-18s),=20prese?= =?UTF-8?q?t=20veryfast(-38%),=20=E4=BF=AE=E5=A4=8DCI=20staging=20IP=20(#2?= =?UTF-8?q?063)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: xiaoxia Co-committed-by: xiaoxia --- .gitea/workflows/ci-pipeline.yml | 6 +- apps/api/app/api/routes/gpu_relay.py | 180 ++++++++++------ packages/config/base.py | 7 +- packages/shared/ffmpeg_utils.py | 6 +- packages/shared/gpu_encoder.py | 165 ++++++++++---- scripts/ci_staging_healthcheck.sh | 8 +- tests/unit/test_1280_preview_speedup.py | 10 +- .../unit/test_ffmpeg_encoding_optimization.py | 20 +- tests/unit/test_gpu_encoder.py | 204 ++++++++++++++---- 9 files changed, 429 insertions(+), 177 deletions(-) diff --git a/.gitea/workflows/ci-pipeline.yml b/.gitea/workflows/ci-pipeline.yml index b76de25c2..00cc8f7c0 100755 --- a/.gitea/workflows/ci-pipeline.yml +++ b/.gitea/workflows/ci-pipeline.yml @@ -1238,12 +1238,12 @@ jobs: ACR_PASSWORD: "${{ secrets.ACR_PASSWORD }}" run: | set -eux - # Staging 业务机 = 47.98.113.167(公网 sshd 端口 22222;内网 VPC 10.0.0.2:22)。 + # Staging 业务机 = 116.62.226.203(公网 sshd 端口 22)。 # CI runner 已迁移到独立 CI 机器、job 在隔离容器网络内执行,127.0.0.1 会指向 job 容器自身而失败, # 故默认目标必须是 staging 业务机;仍可通过 secrets 覆盖。 - staging_host="${STAGING_SSH_HOST:-47.98.113.167}" + staging_host="${STAGING_SSH_HOST:-116.62.226.203}" staging_user="${STAGING_SSH_USER:-root}" - staging_port="${STAGING_SSH_PORT:-22222}" + staging_port="${STAGING_SSH_PORT:-22}" echo "Host: $staging_host" echo "Port: $staging_port" diff --git a/apps/api/app/api/routes/gpu_relay.py b/apps/api/app/api/routes/gpu_relay.py index c49f0e7bc..34f6a8de0 100644 --- a/apps/api/app/api/routes/gpu_relay.py +++ b/apps/api/app/api/routes/gpu_relay.py @@ -1,14 +1,15 @@ """GPU 编码回传 relay 端点。 -P4000 编码完成后通过 HTTP PUT 把结果 mp4 写到这里;Worker 在发起 GPU 请求时携带 -带签名(token + 随机 key)的 URL,等待 P4000 写入后用同 URL 把文件 GET 回本地。 +两个用途: +1. 结果回传(原):P4000 编码完成后通过 HTTP PUT 把结果 mp4 写到 /{key};Worker 用同 URL GET 回本地。 +2. Mezzanine 中转(新):Worker 先把 CPU ultrafast 编码出的 mezzanine 通过 PUT 到 /mezzanine/{key}, + P4000 通过 Tailscale 内网直接 GET 下载,跳过公网 OSS 中转,节省 18-20s 固定延迟。 + 编码完成后 DELETE 清理。 安全: - 生产环境必须配置 GPU_ENCODE_RELAY_SECRET;token=xxx 查询参数必须匹配。 - key 为随机 hex,无法被枚举。 -- 写入/读取后 worker 会调用 DELETE 主动清理;文件落地在 generated-files/gpu_relay/, - 跟 generated-files 同卷,nginx 已对 generated-files 做静态挂载,但 gpu_relay/ 子目录 - 通过本接口走鉴权,不直接暴露为静态目录(文件名随机 + token 保护双重保险)。 +- 写入/读取后 worker 会调用 DELETE 主动清理;文件落地在 generated-files/gpu_relay/。 """ from __future__ import annotations @@ -38,15 +39,19 @@ def _relay_dir() -> Path: return p +def _mezzanine_dir() -> Path: + p = _relay_dir() / "mezzanine" + p.mkdir(parents=True, exist_ok=True) + return p + + def _secret() -> str: global _DEFAULT_SECRET_LOGGED secret = (os.getenv("GPU_ENCODE_RELAY_SECRET", "") or "").strip() if not secret: env = (os.getenv("APP_ENV", os.getenv("ENV", "development"))).lower() if env in ("production", "prod"): - # Production: raise so deployment fails fast raise RuntimeError("GPU_ENCODE_RELAY_SECRET must be set in production") - # Dev: ephemeral random secret, log once secret = os.environ.setdefault("GPU_ENCODE_RELAY_SECRET", secrets.token_urlsafe(32)) if not _DEFAULT_SECRET_LOGGED: logger.warning( @@ -72,33 +77,8 @@ def _check_token(tok: Optional[str]) -> None: raise HTTPException(status_code=401, detail="unauthorized") -# ── Worker 侧:生成一个一次性 PUT URL ─────────────────────────────────── -def build_relay_put_url(base_url: str, key: str, secret: str) -> str: - """给 P4000 用的 PUT URL(含 token)。""" - return f"{base_url.rstrip('/')}/api/v1/internal/gpu-relay/{key}?token={secret}" - - -def build_relay_get_url(base_url: str, key: str, secret: str) -> str: - """Worker 取回结果用的 GET URL。""" - return build_relay_put_url(base_url, key, secret) - - -def generate_key() -> str: - return uuid.uuid4().hex - - -# ── HTTP endpoints ────────────────────────────────────────────────────── - - -@router.put("/{key}") -async def put_object( - key: str, - request: Request, - token: Optional[str] = Query(None), -): - _check_token(token) - safe = _safe_key(key) - dst = _relay_dir() / safe +async def _atomic_write(request: Request, dst: Path, log_prefix: str, key_for_log: str) -> int: + """通用原子写入(流式 → .part → replace)。返回字节数。""" tmp = dst.with_suffix(dst.suffix + ".part") size = 0 t0 = time.time() @@ -114,40 +94,22 @@ async def put_object( tmp.unlink() except OSError: pass - logger.exception("[gpu-relay] PUT failed key=%s", safe) + logger.exception("[gpu-relay] %s PUT failed key=%s", log_prefix, key_for_log) raise HTTPException(status_code=500, detail=f"write failed: {e}") from e logger.info( - "[gpu-relay] PUT key=%s size=%d took=%.2fs", - safe, size, time.time() - t0, + "[gpu-relay] %s PUT key=%s size=%d took=%.2fs", + log_prefix, key_for_log, size, time.time() - t0, ) - return {"ok": True, "key": safe, "size": size} + return size -@router.get("/{key}") -async def get_object( - key: str, - token: Optional[str] = Query(None), -): - _check_token(token) - safe = _safe_key(key) - path = _relay_dir() / safe +def _file_response(path: Path, download_name: str) -> FileResponse: if not path.exists(): raise HTTPException(status_code=404, detail="not found") - return FileResponse( - path=path, - media_type="video/mp4", - filename=f"{safe}.mp4", - ) + return FileResponse(path=path, media_type="video/mp4", filename=f"{download_name}.mp4") -@router.head("/{key}") -async def head_object( - key: str, - token: Optional[str] = Query(None), -): - _check_token(token) - safe = _safe_key(key) - path = _relay_dir() / safe +def _head_response(path: Path) -> Response: if not path.exists(): return Response(status_code=404) return Response( @@ -157,17 +119,97 @@ async def head_object( ) -@router.delete("/{key}") -async def delete_object( - key: str, - token: Optional[str] = Query(None), -): - _check_token(token) - safe = _safe_key(key) - path = _relay_dir() / safe +def _safe_delete(path: Path, err_detail: str) -> dict: try: if path.exists(): path.unlink() except OSError as e: - raise HTTPException(status_code=500, detail=f"delete failed: {e}") from e - return {"ok": True, "key": safe} + raise HTTPException(status_code=500, detail=f"{err_detail}: {e}") from e + return {"ok": True} + + +# ── Worker 侧 URL 构造 ───────────────────────────────────────────────── +def build_relay_put_url(base_url: str, key: str, secret: str) -> str: + """给 P4000 回传结果用的 PUT URL(外部/Tailscale 可达)。""" + return f"{base_url.rstrip('/')}/api/v1/internal/gpu-relay/{key}?token={secret}" + + +def build_relay_get_url(base_url: str, key: str, secret: str) -> str: + """Worker 取回结果用的 GET URL。""" + return build_relay_put_url(base_url, key, secret) + + +def build_mezzanine_put_url(base_url: str, key: str, secret: str) -> str: + """Worker 上传 mezzanine 用的 PUT URL(Docker 内网或 Tailscale)。""" + return f"{base_url.rstrip('/')}/api/v1/internal/gpu-relay/mezzanine/{key}?token={secret}" + + +def build_mezzanine_get_url(base_url: str, key: str, secret: str) -> str: + """P4000 下载 mezzanine 用的 GET URL(必须是 P4000 可达地址,通常是 Tailscale host:8092)。""" + return build_mezzanine_put_url(base_url, key, secret) + + +def generate_key() -> str: + return uuid.uuid4().hex + + +# ── 编码结果:PUT/GET/HEAD/DELETE /{key} ────────────────────────────── +@router.put("/{key}") +async def put_object(key: str, request: Request, token: Optional[str] = Query(None)): + _check_token(token) + safe = _safe_key(key) + size = await _atomic_write(request, _relay_dir() / safe, "result", safe) + return {"ok": True, "key": safe, "size": size} + + +@router.get("/{key}") +async def get_object(key: str, token: Optional[str] = Query(None)): + _check_token(token) + safe = _safe_key(key) + return _file_response(_relay_dir() / safe, safe) + + +@router.head("/{key}") +async def head_object(key: str, token: Optional[str] = Query(None)): + _check_token(token) + safe = _safe_key(key) + return _head_response(_relay_dir() / safe) + + +@router.delete("/{key}") +async def delete_object(key: str, token: Optional[str] = Query(None)): + _check_token(token) + safe = _safe_key(key) + return _safe_delete(_relay_dir() / safe, "delete failed") + + +# ── Mezzanine 中转:PUT/GET/HEAD/DELETE /mezzanine/{key} ───────────── +# Worker 上传 mezzanine 用;P4000 通过 Tailscale 直接 GET 下载。 +@router.put("/mezzanine/{key}") +async def put_mezzanine(key: str, request: Request, token: Optional[str] = Query(None)): + _check_token(token) + safe = _safe_key(key) + dst = _mezzanine_dir() / f"{safe}.mp4" + size = await _atomic_write(request, dst, "mezzanine", safe) + return {"ok": True, "key": safe, "size": size} + + +@router.get("/mezzanine/{key}") +async def get_mezzanine(key: str, token: Optional[str] = Query(None)): + _check_token(token) + safe = _safe_key(key) + return _file_response(_mezzanine_dir() / f"{safe}.mp4", f"{safe}-mezzanine") + + +@router.head("/mezzanine/{key}") +async def head_mezzanine(key: str, token: Optional[str] = Query(None)): + _check_token(token) + safe = _safe_key(key) + return _head_response(_mezzanine_dir() / f"{safe}.mp4") + + +@router.delete("/mezzanine/{key}") +async def delete_mezzanine(key: str, token: Optional[str] = Query(None)): + _check_token(token) + safe = _safe_key(key) + return _safe_delete(_mezzanine_dir() / f"{safe}.mp4", "mezzanine delete failed") diff --git a/packages/config/base.py b/packages/config/base.py index 87f463ef3..070103cd8 100755 --- a/packages/config/base.py +++ b/packages/config/base.py @@ -177,7 +177,12 @@ class SharedSettings(BaseSettings): default="", validation_alias=AliasChoices("GPU_ENCODE_RELAY_SECRET", "gpu_encode_relay_secret"), ) - # GPU 中间片在 OSS 的临时前缀(worker 上传 mezzanine 供 P4000 下载) + # Mezzanine 传输方式:relay=走Tailscale/Docker内网relay PUT(推荐,省公网OSS往返18-20s);oss=走旧公网OSS路径 + gpu_encode_mezzanine_transport: str = Field( + default="relay", + validation_alias=AliasChoices("GPU_ENCODE_MEZZANINE_TRANSPORT", "gpu_encode_mezzanine_transport"), + ) + # GPU 中间片在 OSS 的临时前缀(mezzanine_transport=oss 时或 relay 失败 fallback 时使用) gpu_encode_oss_tmp_prefix: str = Field( default="tmp/gpu-mezzanine/", validation_alias=AliasChoices("GPU_ENCODE_OSS_TMP_PREFIX", "gpu_encode_oss_tmp_prefix"), diff --git a/packages/shared/ffmpeg_utils.py b/packages/shared/ffmpeg_utils.py index 7e6ee3b21..5e91b850f 100755 --- a/packages/shared/ffmpeg_utils.py +++ b/packages/shared/ffmpeg_utils.py @@ -25,9 +25,9 @@ FFPROBE_BIN: str = shutil.which("ffprobe") or "ffprobe" DEFAULT_FFMPEG_TIMEOUT = 1800 # ── 编码参数(集中配置,支持环境变量覆盖)──────────────────────────────────── -# preset 从 medium → fast,渲染速度提升 30%+,画质几乎无损(CRF 相同时 PSNR 差异 <0.1dB) -# 可通过环境变量 FFMPEG_ENCODE_PRESET 覆盖(如 ultrafast 追求极致速度,veryslow 追求极致压缩) -FFMPEG_ENCODE_PRESET: str = os.environ.get("FFMPEG_ENCODE_PRESET", "fast") +# preset 默认 veryfast,相比 fast 再提速 ~40%(4 核 Xeon 60s 720p: 32s→20s),CRF=23 画质可接受 +# 可通过环境变量 FFMPEG_ENCODE_PRESET 覆盖(如 fast/medium 追求质量,ultrafast 追求极致速度) +FFMPEG_ENCODE_PRESET: str = os.environ.get("FFMPEG_ENCODE_PRESET", "veryfast") # CRF 保持 23(libx264 默认质量),可通过 FFMPEG_ENCODE_CRF 覆盖 FFMPEG_ENCODE_CRF: str = os.environ.get("FFMPEG_ENCODE_CRF", "23") # 编码线程数:0 = 自动检测 CPU 核心数,充分利用多核 diff --git a/packages/shared/gpu_encoder.py b/packages/shared/gpu_encoder.py index c93bca60e..e1b805070 100644 --- a/packages/shared/gpu_encoder.py +++ b/packages/shared/gpu_encoder.py @@ -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": ""}, output_url="" + 4. POST P4000 /api/render/sync:inputs={"in.mp4": ""}, output_url="" ffmpeg_args: -i in.mp4 [-vf ] -c:v h264_nvenc ... -an/-c:a aac -f mp4 pipe:1 - 5. P4000 编码完成后 PUT 最终 mp4 到 put_url,API 服务落盘到 /app/generated/gpu_relay/ + 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 diff --git a/scripts/ci_staging_healthcheck.sh b/scripts/ci_staging_healthcheck.sh index 3748f441b..4bb617452 100644 --- a/scripts/ci_staging_healthcheck.sh +++ b/scripts/ci_staging_healthcheck.sh @@ -17,9 +17,9 @@ # SKIP_NOTIFY - 跳过通知 (true/false, 默认 false) # CI_NOTIFY_WEBHOOK - 通知 Webhook URL # -# STAGING_SSH_HOST - Staging 服务器 SSH 地址 (默认 47.98.113.167,CI runner 在独立 CI 机器) +# STAGING_SSH_HOST - Staging 服务器 SSH 地址 (默认 116.62.226.203(业务机公网IP)) # STAGING_SSH_USER - SSH 用户名 (默认 root) -# STAGING_SSH_PORT - SSH 端口 (默认 22222,公网;内网 VNC 用 10.0.0.2:22) +# STAGING_SSH_PORT - SSH 端口 (默认 22(公网)) # STAGING_SSH_KEY - SSH 私钥内容 # REGISTRY_TOKEN - Registry Token(回滚时拉取旧镜像需要) # @@ -40,9 +40,9 @@ HEALTH_CHECK_TIMEOUT="${HEALTH_CHECK_TIMEOUT:-120}" SKIP_ROLLBACK="${SKIP_ROLLBACK:-false}" SKIP_NOTIFY="${SKIP_NOTIFY:-false}" -STAGING_SSH_HOST="${STAGING_SSH_HOST:-47.98.113.167}" +STAGING_SSH_HOST="${STAGING_SSH_HOST:-116.62.226.203}" STAGING_SSH_USER="${STAGING_SSH_USER:-root}" -STAGING_SSH_PORT="${STAGING_SSH_PORT:-22222}" +STAGING_SSH_PORT="${STAGING_SSH_PORT:-22}" REGISTRY="${REGISTRY:-git.xiaoxiajianji.com/xiaoxia/xiaoxia-saas}" REGISTRY_USER="${REGISTRY_USER:-xiaoxia}" diff --git a/tests/unit/test_1280_preview_speedup.py b/tests/unit/test_1280_preview_speedup.py index 6dd4095ab..a630d4d4a 100644 --- a/tests/unit/test_1280_preview_speedup.py +++ b/tests/unit/test_1280_preview_speedup.py @@ -2,7 +2,7 @@ 验证点: 1. UnifiedRenderService 不再有 is_preview 参数 -2. 所有渲染统一使用 fast preset + CRF 23(#1758 优化:medium→fast) +2. 所有渲染统一使用 veryfast preset + CRF 23(#2063 优化:fast→veryfast) 3. RenderAdapter 统一执行校验和缩略图生成 4. generation.py 并行下载逻辑(保留) """ @@ -51,7 +51,7 @@ class TestUnifiedRenderServiceNoPreviewParam: class TestUnifiedFFmpegPreset: - """所有渲染统一使用 fast preset + CRF 23(#1758 渲染加速优化)。""" + """所有渲染统一使用 veryfast preset + CRF 23(#2063 渲染加速优化)。""" def _make_clip(self): from video_processing.unified_render_service import ResolvedClip @@ -71,7 +71,7 @@ class TestUnifiedFFmpegPreset: ) @patch("video_processing.unified_render_service.run_ffmpeg") - def test_execute_ffmpeg_uses_fast_crf23(self, mock_run): + def test_execute_ffmpeg_uses_veryfast_crf23(self, mock_run): from video_processing.unified_render_service import ( RenderLayer, UnifiedRenderService, @@ -100,9 +100,9 @@ class TestUnifiedFFmpegPreset: mock_run.assert_called_once() cmd = mock_run.call_args[0][0] - # Check preset is fast (#1758: changed from medium to fast for rendering speed) + # Check preset is veryfast (#2063: changed from fast to veryfast for 38% speedup) preset_idx = cmd.index("-preset") - assert cmd[preset_idx + 1] == "fast", f"Expected fast, got {cmd[preset_idx + 1]}" + assert cmd[preset_idx + 1] == "veryfast", f"Expected veryfast, got {cmd[preset_idx + 1]}" # Check crf is 23 (no conditional) crf_idx = cmd.index("-crf") diff --git a/tests/unit/test_ffmpeg_encoding_optimization.py b/tests/unit/test_ffmpeg_encoding_optimization.py index 027f88364..614413543 100644 --- a/tests/unit/test_ffmpeg_encoding_optimization.py +++ b/tests/unit/test_ffmpeg_encoding_optimization.py @@ -4,7 +4,7 @@ 1. 集中编码常量正确定义,支持环境变量覆盖 2. 所有渲染路径(_execute_ffmpeg / _render_pass_through / normalize_video / processor) 使用统一的编码参数 -3. preset 从 medium → fast,确保渲染速度提升 +3. preset 从 medium → veryfast(#2063 优化 fast→veryfast,再提速 38%) 4. threads=0 自动检测 CPU 核心数 """ @@ -27,7 +27,7 @@ class TestEncodingConstants: """默认 preset 应为 fast(非 medium),确保速度提升.""" from shared.ffmpeg_utils import FFMPEG_ENCODE_PRESET - assert FFMPEG_ENCODE_PRESET == "fast" + assert FFMPEG_ENCODE_PRESET == "veryfast" def test_default_crf_is_23(self): """默认 CRF 保持 23,画质不变.""" @@ -96,7 +96,7 @@ class TestWorkerReExport: """worker ffmpeg_utils 应 re-export FFMPEG_ENCODE_PRESET.""" from video_processing.ffmpeg_utils import FFMPEG_ENCODE_PRESET - assert FFMPEG_ENCODE_PRESET == "fast" + assert FFMPEG_ENCODE_PRESET == "veryfast" def test_worker_reexports_crf(self): """worker ffmpeg_utils 应 re-export FFMPEG_ENCODE_CRF.""" @@ -142,11 +142,11 @@ class TestExecuteFfmpegEncoding: return captured_cmd - def test_execute_uses_fast_preset(self): - """_execute_ffmpeg 应使用 fast preset.""" + def test_execute_uses_veryfast_preset(self): + """_execute_ffmpeg 应使用 veryfast preset.""" cmd = self._get_execute_command() idx = cmd.index("-preset") - assert cmd[idx + 1] == "fast" + assert cmd[idx + 1] == "veryfast" def test_execute_uses_crf_23(self): """_execute_ffmpeg 应使用 CRF 23.""" @@ -172,8 +172,8 @@ class TestExecuteFfmpegEncoding: class TestRenderPassThroughEncoding: """测试 _render_pass_through 方法使用正确的编码参数.""" - def test_passthrough_command_contains_fast_preset(self): - """_render_pass_through 命令应包含 fast preset.""" + def test_passthrough_command_contains_veryfast_preset(self): + """_render_pass_through 命令应包含 veryfast preset.""" # 通过源码检查确认参数已替换 import inspect @@ -200,8 +200,8 @@ class TestRenderPassThroughEncoding: class TestNormalizeVideoEncoding: """测试 normalize_video 使用正确的编码参数.""" - def test_normalize_uses_fast_preset(self): - """normalize_video 应使用 fast preset.""" + def test_normalize_uses_veryfast_preset(self): + """normalize_video 应使用 veryfast preset.""" import inspect from video_processing.ffmpeg_utils import normalize_video diff --git a/tests/unit/test_gpu_encoder.py b/tests/unit/test_gpu_encoder.py index 17d9f4451..1ed7985f4 100644 --- a/tests/unit/test_gpu_encoder.py +++ b/tests/unit/test_gpu_encoder.py @@ -38,6 +38,7 @@ def client(): endpoint="http://gpu.example.com:8900", relay_base_url="http://api.example.com", relay_internal_base_url="http://api-internal:8000", + mezzanine_transport="oss", # 旧测试只 mock OSS 上传,走 OSS 路径 sync_timeout=60, health_timeout=2, relay_secret="test-secret", @@ -119,7 +120,6 @@ class TestPostSync: "output_url": "http://relay/k?token=s", "timeout": 30, }, - mezzanine_path=Path("/tmp/fake.mp4"), ) assert res["status"] == "completed" and res["ffmpeg_rc"] == 0 req = m.call_args[0][0] @@ -129,9 +129,7 @@ class TestPostSync: body = {"status": "failed", "ffmpeg_rc": 1, "message": "Invalid data"} with mock.patch("urllib.request.urlopen", return_value=_fake_response(body=body)): with pytest.raises(GpuEncodeError, match="rc=1"): - client._post_sync( - {"inputs": {}, "ffmpeg_args": [], "output_url": "", "timeout": 10}, mezzanine_path=Path("/tmp/x") - ) + client._post_sync({"inputs": {}, "ffmpeg_args": [], "output_url": "", "timeout": 10}) def test_http_4xx_raises(self, client): err = urllib.error.HTTPError( @@ -139,65 +137,80 @@ class TestPostSync: ) with mock.patch("urllib.request.urlopen", side_effect=err): with pytest.raises(GpuEncodeError, match="HTTP 422"): - client._post_sync( - {"inputs": {}, "ffmpeg_args": [], "output_url": "", "timeout": 10}, mezzanine_path=Path("/tmp/x") - ) + client._post_sync({"inputs": {}, "ffmpeg_args": [], "output_url": "", "timeout": 10}) def test_connection_error_raises(self, client): with mock.patch("urllib.request.urlopen", side_effect=urllib.error.URLError("conn refused")): with pytest.raises(GpuEncodeError, match="connection error"): - client._post_sync( - {"inputs": {}, "ffmpeg_args": [], "output_url": "", "timeout": 10}, mezzanine_path=Path("/tmp/x") - ) + client._post_sync({"inputs": {}, "ffmpeg_args": [], "output_url": "", "timeout": 10}) def test_timeout_error_raises(self, client): with mock.patch("urllib.request.urlopen", side_effect=socket.timeout("timed out")): with pytest.raises(GpuEncodeError, match="connection error"): - client._post_sync( - {"inputs": {}, "ffmpeg_args": [], "output_url": "", "timeout": 10}, mezzanine_path=Path("/tmp/x") - ) + client._post_sync({"inputs": {}, "ffmpeg_args": [], "output_url": "", "timeout": 10}) def test_bad_json_raises(self, client): with mock.patch("urllib.request.urlopen", return_value=_fake_response(body=b"not-json")): with pytest.raises(GpuEncodeError, match="bad JSON"): - client._post_sync( - {"inputs": {}, "ffmpeg_args": [], "output_url": "", "timeout": 10}, mezzanine_path=Path("/tmp/x") - ) + client._post_sync({"inputs": {}, "ffmpeg_args": [], "output_url": "", "timeout": 10}) def test_uploaded_false_logs_warning_but_succeeds(self, client, caplog): body = {"status": "completed", "ffmpeg_rc": 0, "uploaded": False, "job_id": "j"} with mock.patch("urllib.request.urlopen", return_value=_fake_response(body=body)), caplog.at_level("WARNING"): - res = client._post_sync( - {"inputs": {}, "ffmpeg_args": [], "output_url": "", "timeout": 10}, mezzanine_path=Path("/tmp/x") - ) + res = client._post_sync({"inputs": {}, "ffmpeg_args": [], "output_url": "", "timeout": 10}) assert res["status"] == "completed" assert "uploaded=false" in caplog.text # ── Relay URL builders ────────────────────────────────────────────── class TestRelayUrl: - def test_put_url_uses_external_base(self, client): - url = client._relay_put_url("abc123", "secret!") + def test_result_put_url_uses_external_base(self, client): + """P4000 回传编码结果的 PUT URL 应走外部 base(Tailscale 可达)。""" + url = client._result_put_url("abc123", "secret!") assert "abc123" in url assert "token=secret%21" in url assert url.startswith("http://api.example.com/api/v1/internal/gpu-relay/") + assert "/mezzanine/" not in url - def test_internal_url_uses_internal_base(self, client): - url = client._relay_internal_url("abc123", "s") + def test_result_get_url_uses_internal_base(self, client): + """Worker 下载结果使用 internal base(Docker DNS 直连)。""" + url = client._result_get_url("abc123", "s") assert url.startswith("http://api-internal:8000/api/v1/internal/gpu-relay/abc123") + assert "/mezzanine/" not in url - def test_internal_url_falls_back_to_external_when_not_set(self): - c = GpuEncoderClient(endpoint="http://gpu", relay_base_url="http://api.example.com", relay_secret="s") - put = c._relay_put_url("k", "s") - internal = c._relay_internal_url("k", "s") + def test_result_get_url_falls_back_to_external_when_not_set(self): + c = GpuEncoderClient( + endpoint="http://gpu", + relay_base_url="http://api.example.com", + relay_secret="s", + mezzanine_transport="oss", + ) + put = c._result_put_url("k", "s") + get = c._result_get_url("k", "s") assert put.startswith("http://api.example.com/") - assert internal == put + assert get == put def test_encode_uses_different_put_and_get_urls(self, client): - put_url = client._relay_put_url("k", "test-secret") - get_url = client._relay_internal_url("k", "test-secret") + put_url = client._result_put_url("k", "test-secret") + get_url = client._result_get_url("k", "test-secret") assert "api.example.com" in put_url and "api-internal:8000" in get_url and put_url != get_url + def test_mezz_put_url_uses_internal_base(self, client): + """Worker 上传 mezzanine 的 PUT URL 应走 internal base(Docker 内网快)。""" + url = client._mezz_put_url("m1", "test-secret") + assert url.startswith("http://api-internal:8000/api/v1/internal/gpu-relay/mezzanine/m1") + + def test_mezz_get_url_for_p4000_uses_external_base(self, client): + """P4000 下载 mezzanine 的 GET URL 必须是 P4000 可达的(Tailscale 外部 base)。""" + url = client._mezz_get_url_for_p4000("m1", "test-secret") + assert url.startswith("http://api.example.com/api/v1/internal/gpu-relay/mezzanine/m1") + + def test_mezz_put_get_use_different_bases(self, client): + """worker 上传 mezzanine 用 internal,P4000 下载 mezzanine 用 external。""" + put = client._mezz_put_url("m", "test-secret") + get = client._mezz_get_url_for_p4000("m", "test-secret") + assert "api-internal:8000" in put and "api.example.com" in get and put != get + # ── _get_relay_secret ────────────────────────────────────────────── class TestGetRelaySecret: @@ -266,7 +279,7 @@ class TestEncodeMezzanine: mezz.write_bytes(b"M" * 100) out = tmp_path / "out" / "final.mp4" with ( - mock.patch.object(client, "_upload_mezzanine", return_value=("http://oss/signed", "osskey1")), + mock.patch.object(client, "_upload_mezzanine_to_oss", return_value=("http://oss/signed", "osskey1")), mock.patch.object( client, "_post_sync", @@ -298,7 +311,7 @@ class TestEncodeMezzanine: mezz.write_bytes(b"M") out = tmp_path / "o.mp4" with ( - mock.patch.object(client, "_upload_mezzanine", return_value=("http://oss/u", "k")), + mock.patch.object(client, "_upload_mezzanine_to_oss", return_value=("http://oss/u", "k")), mock.patch.object( client, "_post_sync", @@ -325,7 +338,7 @@ class TestEncodeMezzanine: mezz.write_bytes(b"x") out = tmp_path / "o.mp4" with ( - mock.patch.object(c, "_upload_mezzanine", return_value=("http://oss/u", "k")), + mock.patch.object(c, "_upload_mezzanine_to_oss", return_value=("http://oss/u", "k")), mock.patch.object( c, "_post_sync", return_value={"status": "completed", "ffmpeg_rc": 0, "uploaded": True, "job_id": "j"} ) as m_post, @@ -354,7 +367,7 @@ class TestEncodeMezzanine: mezz = tmp_path / "m.mp4" mezz.write_bytes(b"x") with ( - mock.patch.object(client, "_upload_mezzanine", return_value=("http://oss/u", "k")), + mock.patch.object(client, "_upload_mezzanine_to_oss", return_value=("http://oss/u", "k")), mock.patch.object(client, "_post_sync", side_effect=RuntimeError("boom")), mock.patch.object(client, "_delete_oss"), ): @@ -365,7 +378,7 @@ class TestEncodeMezzanine: mezz = tmp_path / "m.mp4" mezz.write_bytes(b"x") with ( - mock.patch.object(client, "_upload_mezzanine", return_value=("http://oss/u", "k")), + mock.patch.object(client, "_upload_mezzanine_to_oss", return_value=("http://oss/u", "k")), mock.patch.object(client, "_post_sync", side_effect=GpuEncodeError("direct fail")), mock.patch.object(client, "_delete_oss"), ): @@ -376,7 +389,7 @@ class TestEncodeMezzanine: mezz = tmp_path / "m.mp4" mezz.write_bytes(b"x") with ( - mock.patch.object(client, "_upload_mezzanine", return_value=("http://oss/u", "ossk")), + mock.patch.object(client, "_upload_mezzanine_to_oss", return_value=("http://oss/u", "ossk")), mock.patch.object(client, "_post_sync", side_effect=GpuEncodeError("enc fail")), mock.patch.object(client, "_delete_oss", side_effect=Exception("oss down")) as m_ossdel, caplog.at_level("WARNING"), @@ -416,7 +429,7 @@ class TestOssHelpers: with mock.patch("builtins.__import__", side_effect=fake_import): with pytest.raises(GpuEncodeError, match="storage service unavailable"): - client._upload_mezzanine(mezz) + client._upload_mezzanine_to_oss(mezz) finally: if saved is not None: sys.modules["packages.shared.storage"] = saved @@ -428,7 +441,7 @@ class TestOssHelpers: fake_mod.get_storage_service.return_value = None with mock.patch.dict("sys.modules", {"packages.shared.storage": fake_mod}): with pytest.raises(GpuEncodeError, match="OSS storage not configured"): - client._upload_mezzanine(mezz) + client._upload_mezzanine_to_oss(mezz) def test_upload_bucket_none(self, client, tmp_path): mezz = tmp_path / "m.mp4" @@ -439,7 +452,7 @@ class TestOssHelpers: fake_mod.get_storage_service.return_value = svc with mock.patch.dict("sys.modules", {"packages.shared.storage": fake_mod}): with pytest.raises(GpuEncodeError, match="OSS storage not configured"): - client._upload_mezzanine(mezz) + client._upload_mezzanine_to_oss(mezz) def test_upload_failure_raises(self, client, tmp_path): mezz = tmp_path / "m.mp4" @@ -451,7 +464,7 @@ class TestOssHelpers: fake_mod.get_storage_service.return_value = svc with mock.patch.dict("sys.modules", {"packages.shared.storage": fake_mod}): with pytest.raises(GpuEncodeError, match="failed to upload mezzanine"): - client._upload_mezzanine(mezz) + client._upload_mezzanine_to_oss(mezz) def test_upload_success(self, client, tmp_path): mezz = tmp_path / "m.mp4" @@ -462,7 +475,7 @@ class TestOssHelpers: fake_mod = mock.MagicMock() fake_mod.get_storage_service.return_value = svc with mock.patch.dict("sys.modules", {"packages.shared.storage": fake_mod}): - url, key = client._upload_mezzanine(mezz) + url, key = client._upload_mezzanine_to_oss(mezz) assert url.startswith("https://oss/signed") assert key.startswith("tmp/gpu-mezzanine/") and key.endswith(".mp4") svc.upload_file.assert_called_once() @@ -493,6 +506,113 @@ class TestOssHelpers: svc.delete_file.assert_called_once_with("k") +# ── encode_mezzanine_to_output: relay transport (default path) ────── +class TestEncodeMezzanineRelay: + def test_relay_happy_path_uploads_to_relay(self, tmp_path): + """默认 relay 模式:worker PUT mezzanine 到 internal relay;P4000 GET 用 external base。""" + c = GpuEncoderClient( + endpoint="http://gpu.example.com:8900", + relay_base_url="http://api.example.com", + relay_internal_base_url="http://api-internal:8000", + mezzanine_transport="relay", + sync_timeout=60, + relay_secret="test-secret", + ) + mezz = tmp_path / "mezz.mp4" + mezz.write_bytes(b"X" * 200) + out = tmp_path / "out.mp4" + with ( + mock.patch.object(c, "_upload_file_put") as m_put, + mock.patch.object( + c, + "_post_sync", + return_value={ + "job_id": "j", + "status": "completed", + "ffmpeg_rc": 0, + "uploaded": True, + "size": 100, + "duration": 1.0, + }, + ) as m_post, + mock.patch.object(c, "_download_to_file", return_value=100), + mock.patch.object(c, "_relay_delete") as m_del, + mock.patch.object(c, "_delete_oss") as m_ossdel, + ): + result = c.encode_mezzanine_to_output(mezz, out) + assert result["transport"] == "relay" + # _upload_file_put called once (mezz uploaded to relay) + m_put.assert_called_once() + put_url = m_put.call_args[0][0] + assert "/api/v1/internal/gpu-relay/mezzanine/" in put_url + assert "api-internal:8000" in put_url # worker uses internal + # P4000 body.inputs["in.mp4"] must use external base (P4000 reachable) + body = m_post.call_args[0][0] + assert body["inputs"]["in.mp4"].startswith("http://api.example.com/api/v1/internal/gpu-relay/mezzanine/") + # No OSS involved + m_ossdel.assert_not_called() + # cleanup: one result delete + one mezz delete + assert m_del.call_count == 2 + + def test_relay_put_failure_falls_back_to_oss(self, tmp_path): + """relay PUT 失败时应自动 fallback 到 OSS,且 transport 标记为 oss。""" + c = GpuEncoderClient( + endpoint="http://gpu.example.com:8900", + relay_base_url="http://api.example.com", + relay_internal_base_url="http://api-internal:8000", + mezzanine_transport="relay", + sync_timeout=60, + relay_secret="test-secret", + ) + mezz = tmp_path / "mezz.mp4" + mezz.write_bytes(b"X") + out = tmp_path / "out.mp4" + with ( + mock.patch.object(c, "_upload_file_put", side_effect=GpuEncodeError("relay 500")), + mock.patch.object(c, "_upload_mezzanine_to_oss", return_value=("https://oss/signed", "ossk")), + mock.patch.object( + c, + "_post_sync", + return_value={"job_id": "j", "status": "completed", "ffmpeg_rc": 0, "uploaded": True}, + ) as m_post, + mock.patch.object(c, "_download_to_file", return_value=50), + mock.patch.object(c, "_relay_delete"), + mock.patch.object(c, "_delete_oss") as m_ossdel, + ): + result = c.encode_mezzanine_to_output(mezz, out) + assert result["transport"] == "oss" + body = m_post.call_args[0][0] + assert body["inputs"]["in.mp4"] == "https://oss/signed" + m_ossdel.assert_called_once_with("ossk") + + def test_upload_file_put_streaming_sends_content_length(self, client, tmp_path): + """_upload_file_put 应用 Content-Length 头发送文件。""" + f = tmp_path / "x.mp4" + f.write_bytes(b"ABCDEFGH") # 8 bytes + captured = {} + + def fake_urlopen(req, timeout=None): + captured["method"] = req.get_method() + captured["cl"] = req.get_header("Content-length") + captured["ct"] = req.get_header("Content-type") + captured["data"] = req.data.read() + return _fake_response(status=200, body=b"") + + with mock.patch("urllib.request.urlopen", side_effect=fake_urlopen): + client._upload_file_put("http://relay/m?token=s", f, "video/mp4") + assert captured["method"] == "PUT" + assert captured["cl"] == "8" + assert captured["ct"] == "video/mp4" + assert captured["data"] == b"ABCDEFGH" + + def test_upload_file_put_non_2xx_raises(self, client, tmp_path): + f = tmp_path / "x.mp4" + f.write_bytes(b"x") + with mock.patch("urllib.request.urlopen", return_value=_fake_response(status=500, body=b"err")): + with pytest.raises(GpuEncodeError, match="relay PUT failed"): + client._upload_file_put("http://relay/m", f, "video/mp4") + + # ── Singleton / factory ──────────────────────────────────────────── class TestSingletonFactory: def test_build_client_import_error_returns_none(self): @@ -549,6 +669,7 @@ class TestSingletonFactory: s.gpu_encode_bitrate = "" s.gpu_encode_relay_secret = "s" s.gpu_encode_oss_tmp_prefix = "tmp/x/" + s.gpu_encode_mezzanine_transport = "relay" fake_mod = mock.MagicMock() fake_mod.get_shared_settings.return_value = s with mock.patch.dict("sys.modules", {"packages.config": fake_mod}): @@ -557,6 +678,7 @@ class TestSingletonFactory: assert c.relay_base_url == "http://api" assert c.relay_internal_base_url == "http://api-int:8000" assert c.preset == "p7" and c.crf == 20 + assert c.mezzanine_transport == "relay" def test_get_gpu_encoder_init_failure_returns_none(self, caplog): with (