perf(gpu): mezzanine relay直传跳过公网OSS(-18s), preset veryfast(-38%), 修复CI staging IP #2063

Merged
auto-approve-bot merged 5 commits from fix/gpu-mezzanine-relay-transport into develop 2026-09-27 14:48:55 +08:00
9 changed files with 429 additions and 177 deletions
+3 -3
View File
@@ -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"
+111 -69
View File
@@ -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")
+6 -1
View File
@@ -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"),
+3 -3
View File
@@ -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 核心数,充分利用多核
+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
+4 -4
View File
@@ -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}"
+5 -5
View File
@@ -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")
+10 -10
View File
@@ -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
+163 -41
View File
@@ -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 (