Files
xiaoxia-saas/packages/shared/gpu_encoder.py
T
xiaoxia 8abdeb9551
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 2s
CI/CD Pipeline / Check push changed paths (pull_request) 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 / 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 / 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 / Validate - Style (pull_request) Has been skipped
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Has been skipped
CI/CD Pipeline / Validate - Security (pull_request) Has been skipped
CI/CD Pipeline / Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / Integration 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 / Check push changed paths (push) Successful in 11s
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 1m1s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 1m9s
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 32s
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 / PR Build Web Image (pull_request) Successful in 1m48s
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m0s
CI/CD Pipeline / CI Gate (pull_request) Successful in 2s
CI/CD Pipeline / Build Staging Web Image (push) Successful in 49s
CI/CD Pipeline / Retag skipped Staging Web Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging API Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Has been skipped
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 53s
CI/CD Pipeline / Integration Tests (push) Successful in 3m10s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m24s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been skipped
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 3m37s
CI/CD Pipeline / Validate - Style (push) Successful in 4m14s
CI/CD Pipeline / ACR Image Cleanup (push) Successful in 2m18s
CI/CD Pipeline / Frontend Unit Tests (push) Failing after 5m56s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 3m29s
AI Code Review / AI Code Review (pull_request) Successful in 7m9s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 4m27s
CI/CD Pipeline / Validate - Security (push) Successful in 8m48s
CI/CD Pipeline / Unit Tests (push) Successful in 10m43s
CI/CD Pipeline / Build Production API Image (push) Has been skipped
CI/CD Pipeline / Build Production Web Image (push) Has been skipped
CI/CD Pipeline / Build Production Worker Image (push) Has been skipped
CI/CD Pipeline / CI Gate (push) Has been skipped
CI/CD Pipeline / Deploy Production (push) Has been skipped
CI/CD Pipeline / Canary Release to Production (push) Has been skipped
CI/CD Pipeline / Production Browser E2E (push) Has been skipped
fix(P0): gpu-encoder 冷启动防护 — pre-warm + 首字节短超时兜底 (#2072)
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com>
Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
2026-09-28 00:08:49 +08:00

570 lines
26 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""P4000 NVENC 远程编码客户端。
完整链路(encode_video_file):
1. CPU 滤镜已在本地生成 mezzanine 中间片(libx264 ultrafast)
2. 通过 HTTP PUT 把 mezzanine 上传到 relay(走 Tailscale/Docker 内网,~1s 完成)
- 失败则 fallback 到 OSS 上传(旧路径,兼容没有 :8092 内网可达的环境)
3. 生成 relay 一次性 key,构造两个带 token 的 URL:
- put_url:给 P4000 回传结果,走 relay_base_url(Tailscale host:8092)
- get/del_url:worker 自己下载+清理用,走 relay_internal_base_url(Docker DNS 直连 API)
4. 【冷启动防护】距上次成功通信 >60s 时,先 GET /health 预热 Tailscale 链路(短超时快速失败)
5. 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
- 首字节用短超时(默认20s),避免链路卡死空等上百秒;首字节到达后放宽到 ffmpeg_timeout+60s
6. P4000 从 relay GET mezzanine → h264_nvenc 编码 → PUT 最终 mp4 到 put_url
7. 本客户端通过 get_url(Docker 内网)下载最终文件到 output_path,然后 DELETE 清理
8. 删除 relay 上的 mezzanine 临时文件(以及 OSS fallback 的 key)
任何环节失败抛 GpuEncodeError,调用方应 fallback 到 CPU libx264。
冷启动/链路卡顿背景(2026-09-27 实测):P4000 与 staging 之间走 Tailscale,长时间空闲
(>7h)后首次请求曾出现 150s 延迟才真正开始下载 mezzanine,期间 ffmpeg 尚未启动、GPU 空闲。
根因在服务端/网络层(可能是 Tailscale DERP 打洞或 httpx 连接池重建),本客户端通过
pre_warm + 首字节短超时做兜底:预热打通链路 + 20s 内收不到首字节就快速失败让 CPU fallback,
不再让用户等满 150s+。
"""
from __future__ import annotations
import json
import logging
import os
import socket
import time
import urllib.error
import urllib.parse
import urllib.request
import uuid
from dataclasses import dataclass
from pathlib import Path
from typing import Any, Optional
logger = logging.getLogger(__name__)
class GpuEncodeError(RuntimeError):
"""GPU 编码失败(网络/超时/ffmpeg/upload/download 任一环节)。调用方应 fallback 到 CPU。"""
@dataclass
class GpuHealth:
healthy: bool
worker: str = ""
gpu_name: str = ""
nvenc_h264: bool = False
nvenc_hevc: bool = False
error: str = ""
@property
def ready(self) -> bool:
return self.healthy and self.nvenc_h264
class GpuEncoderClient:
def __init__(
self,
endpoint: str,
relay_base_url: str,
*,
relay_internal_base_url: str = "",
# Mezzanine 上传:默认走 relay(Tailscale/Docker 内网);设为 "oss" 强制走旧 OSS 路径
mezzanine_transport: str = "relay",
sync_timeout: int = 300,
health_timeout: float = 3.0,
# 提交编码任务前先发一次 /health 预热 Tailscale 链路,避免长时间空闲后首次请求
# 因 DERP 打洞/NAT 映射过期/Tailscale 连接重建而阻塞上百秒。
pre_warm: bool = True,
# POST 首次响应超时:P4000 已收到请求后应该在数秒内开始下载 inputs;
# 如果超过这个值还没收到任何响应字节,说明链路/服务卡住,快速失败让调用方 fallback CPU。
# 注意:ffmpeg 编码本身靠 body.timeout 控制(300s),不应该被这个超时影响。
post_first_byte_timeout: float = 20.0,
vcodec: str = "h264_nvenc",
preset: str = "p4",
crf: int = 23,
bitrate: str = "",
relay_secret: str = "",
oss_tmp_prefix: str = "tmp/gpu-mezzanine/",
) -> None:
self.endpoint = endpoint.rstrip("/")
self.relay_base_url = relay_base_url.rstrip("/")
# Worker→API 内网访问地址(Docker DNS 直连,如 http://xiaoxia-api-staging:8000)。
# 未配置时回退到 relay_base_url(本地开发/单节点)。
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.pre_warm = pre_warm
self.post_first_byte_timeout = post_first_byte_timeout
self.vcodec = vcodec
self.preset = preset
self.crf = crf
self.bitrate = bitrate
self._relay_secret = relay_secret
self.oss_tmp_prefix = oss_tmp_prefix.rstrip("/") + "/" if oss_tmp_prefix else "tmp/gpu-mezzanine/"
# 上次与 P4000 成功通信的时间戳(用于判断是否需要 pre_warm 预热)
self._last_ok_ts: float = 0.0
RELAY_PATH_PREFIX = "/api/v1/internal/gpu-relay"
# ------------------------------------------------------------------
# URL builders
# ------------------------------------------------------------------
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_result_url(self, base_url: str, key: str, secret: str) -> str:
return self._relay_url_from_base(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
# ------------------------------------------------------------------
def check_health(self) -> GpuHealth:
url = f"{self.endpoint}/health"
try:
with urllib.request.urlopen(url, timeout=self.health_timeout) as resp:
data = json.loads(resp.read().decode("utf-8"))
except (urllib.error.URLError, socket.timeout, TimeoutError, json.JSONDecodeError, ConnectionError) as e:
return GpuHealth(healthy=False, error=f"health probe failed: {e}")
try:
h = GpuHealth(
healthy=data.get("status") == "healthy",
worker=str(data.get("worker", "")),
gpu_name=(data.get("gpu") or {}).get("name", ""),
nvenc_h264=bool((data.get("nvenc") or {}).get("h264_nvenc")),
nvenc_hevc=bool((data.get("nvenc") or {}).get("hevc_nvenc")),
)
if h.healthy:
self._last_ok_ts = time.time()
return h
except Exception as e: # noqa: BLE001
return GpuHealth(healthy=False, error=f"malformed health response: {e}")
# ------------------------------------------------------------------
# High-level: encode a mezzanine file to final output
# ------------------------------------------------------------------
def encode_mezzanine_to_output(
self,
mezzanine_path: Path,
output_path: Path,
*,
extra_video_args: Optional[list[str]] = None,
audio_args: Optional[list[str]] = None,
timeout: Optional[int] = None,
) -> dict[str, Any]:
"""把 mezzanine(CPU 滤镜已完成)交给 P4000 NVENC 编码,结果写到 output_path。
传输:默认通过 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}")
if not self.relay_base_url:
raise GpuEncodeError("gpu_encode_relay_base_url not configured")
timeout = timeout or self.sync_timeout
t_total = time.time()
oss_key: Optional[str] = None
mezz_key: Optional[str] = None
result_key: Optional[str] = None
input_url: str = ""
used_transport = self.mezzanine_transport
try:
secret = self._get_relay_secret()
# 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"]
if extra_video_args:
ffmpeg_args.extend(extra_video_args)
ffmpeg_args.extend(["-c:v", self.vcodec, "-preset", self.preset])
if self.bitrate:
ffmpeg_args.extend(["-b:v", self.bitrate])
else:
ffmpeg_args.extend(["-cq", str(self.crf)])
ffmpeg_args.extend(["-pix_fmt", "yuv420p", "-movflags", "+faststart"])
if audio_args:
ffmpeg_args.extend(audio_args)
else:
ffmpeg_args.append("-an")
ffmpeg_args.extend(["-f", "mp4", "pipe:1"])
# 4. pre-warm then call P4000 sync render
self._warm_up_if_needed()
body = {
"inputs": {"in.mp4": input_url},
"ffmpeg_args": ffmpeg_args,
"output_url": put_url,
"timeout": int(timeout),
}
job = self._post_sync(body)
self._last_ok_ts = time.time()
logger.info(
"[gpu-encoder] P4000 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 result
self._relay_delete(del_result_url)
logger.info(
"[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),
"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)
# ------------------------------------------------------------------
# Internal helpers
# ------------------------------------------------------------------
def _get_relay_secret(self) -> str:
if self._relay_secret:
return self._relay_secret
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")
raise GpuEncodeError("GPU_ENCODE_RELAY_SECRET not set")
return secret
def _warm_up_if_needed(self) -> None:
"""POST 前预热:如果距上次成功通信超过 idle 阈值,先打 /health 打通 Tailscale 链路。
背景:Tailscale 在长时间空闲(几小时)后,到对端的直连 NAT 映射可能过期,
首次请求会走 DERP 中继打洞;极少数情况下打洞/重连会卡住上百秒(曾观测到 150s 延迟)。
预热请求本身走短超时快速失败,不会阻塞主流程;预热成功后再发 POST。
"""
if not self.pre_warm:
return
idle = time.time() - self._last_ok_ts
# 空闲超过 60s 才预热(正常流水线里相邻任务间隔通常 <10s,没必要每次都打)
if idle < 60:
return
url = f"{self.endpoint}/health"
t0 = time.time()
try:
with urllib.request.urlopen(url, timeout=min(self.health_timeout, 3.0)) as resp:
resp.read()
self._last_ok_ts = time.time()
logger.debug("[gpu-encoder] pre-warm ok: took=%.2fs idle=%.0fs", time.time() - t0, idle)
except (urllib.error.URLError, socket.timeout, TimeoutError, ConnectionError, OSError) as e:
# 预热失败不致命——主 POST 会带自己的超时,再失败就抛 GpuEncodeError 让调用方 fallback
logger.warning("[gpu-encoder] pre-warm probe failed (will try POST anyway): %s", e)
def _post_sync(self, body: dict[str, Any]) -> dict[str, Any]:
url = f"{self.endpoint}/api/render/sync"
ffmpeg_timeout = body.get("timeout", self.sync_timeout)
# 连接 + 首字节用短超时(防链路卡死数百秒);首字节到达后给 ffmpeg 留足编码+上传时间
# Python urllib 的 timeout 是整个请求总超时,所以用"两段式":
# 阶段1:先 read(1) 拿首字节,用短超时;
# 阶段2:再 read() 读完整 body,用 ffmpeg_timeout+60。
connect_timeout = min(max(self.post_first_byte_timeout, 5.0), 30.0)
payload = json.dumps(body).encode("utf-8")
req = urllib.request.Request(
url,
data=payload,
headers={"Content-Type": "application/json"},
method="POST",
)
t0 = time.time()
first_byte_ok = False
resp = None
try:
resp = urllib.request.urlopen(req, timeout=connect_timeout)
# 读首字节 —— 如果 P4000/链路卡死,这里会在 connect_timeout 内抛超时
first_chunk = resp.read(1)
first_byte_ok = True
logger.debug(
"[gpu-encoder] P4000 first byte in %.2fs (connect_timeout=%.1fs)",
time.time() - t0,
connect_timeout,
)
# 剩余用长超时(给底层socket放宽时限;如果是mock/不支持,则跳过)
try:
resp.fp._sock.settimeout(ffmpeg_timeout + 60)
except (AttributeError, OSError):
pass
rest = resp.read()
raw = (first_chunk + rest).decode("utf-8")
resp.close()
resp = None
except urllib.error.HTTPError as e:
detail = e.read().decode("utf-8", errors="replace")[:1000]
raise GpuEncodeError(f"P4000 HTTP {e.code}: {detail}") from e
except (urllib.error.URLError, socket.timeout, TimeoutError, ConnectionError, OSError) as e:
waited = time.time() - t0
hint = "first-byte" if not first_byte_ok else "ffmpeg/upload"
# 统一以 "connection error" 开头,便于上层 fallback 逻辑用关键词识别;
# 末尾再附带具体错误(timed out / refused ...)供排障
raise GpuEncodeError(
f"P4000 {hint} connection error after {waited:.1f}s "
f"(connect_timeout={connect_timeout:.0f}s, ffmpeg_timeout={ffmpeg_timeout}s): {e}"
) from e
finally:
if resp is not None:
try:
resp.close()
except Exception:
pass
try:
result = json.loads(raw)
except json.JSONDecodeError as e:
raise GpuEncodeError(f"P4000 bad JSON: {raw[:500]}") from e
dt = time.time() - t0
status = result.get("status")
ffmpeg_rc = result.get("ffmpeg_rc")
uploaded = result.get("uploaded")
if status != "completed" or ffmpeg_rc != 0:
err = result.get("message") or result.get("error") or "unknown"
raise GpuEncodeError(f"P4000 job failed: status={status} rc={ffmpeg_rc} err={err!s:.500}")
if not uploaded:
logger.warning("[gpu-encoder] P4000 reports uploaded=false (will verify via download)")
result["_roundtrip"] = dt
return result
def _download_to_file(self, url: str, output_path: Path) -> int:
"""GET url → write to output_path. Returns bytes written."""
tmp = output_path.with_suffix(output_path.suffix + ".gpu_tmp")
size = 0
try:
with urllib.request.urlopen(url, timeout=self.sync_timeout) as resp:
if resp.status != 200:
raise GpuEncodeError(f"relay GET returned HTTP {resp.status}")
with open(tmp, "wb") as f:
while True:
chunk = resp.read(1024 * 256)
if not chunk:
break
f.write(chunk)
size += len(chunk)
if size == 0:
raise GpuEncodeError("relay returned empty file")
os.replace(tmp, output_path)
return size
except (urllib.error.URLError, socket.timeout, TimeoutError, ConnectionError) as e:
if tmp.exists():
try:
tmp.unlink()
except OSError:
pass
raise GpuEncodeError(f"failed to download from relay: {e}") from e
def _relay_delete(self, url: str) -> None:
try:
req = urllib.request.Request(url, method="DELETE")
with urllib.request.urlopen(req, timeout=10) as resp:
resp.read()
except Exception as e: # noqa: BLE001
logger.debug("[gpu-encoder] relay cleanup delete failed: %s", e)
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 (fallback)
# ------------------------------------------------------------------
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
except ImportError as e:
raise GpuEncodeError(f"storage service unavailable: {e}") from e
storage = get_storage_service()
if storage is None or storage.bucket is None:
raise GpuEncodeError("OSS storage not configured; cannot upload mezzanine")
key = f"{self.oss_tmp_prefix}{uuid.uuid4().hex}.mp4"
try:
storage.upload_file(str(path), key, content_type="video/mp4")
except Exception as e: # noqa: BLE001
raise GpuEncodeError(f"failed to upload mezzanine to OSS: {e}") from e
signed = storage.get_download_url(key, expires_seconds=3600)
return signed, key
def _delete_oss(self, key: str) -> None:
try:
from packages.shared.storage import get_storage_service
storage = get_storage_service()
if storage is not None and storage.bucket is not None:
storage.delete_file(key)
except Exception as e: # noqa: BLE001
logger.debug("[gpu-encoder] OSS delete %s failed: %s", key, e)
# ── Singleton factory ────────────────────────────────────────────────────
_default_client: Optional[GpuEncoderClient] = None
_default_client_initialized: bool = False
def _build_client_from_settings() -> Optional[GpuEncoderClient]:
try:
from packages.config import get_shared_settings
settings = get_shared_settings()
except Exception: # noqa: BLE001
return None
if not getattr(settings, "enable_gpu_encode", False):
return None
endpoint = (getattr(settings, "gpu_encode_endpoint", "") or "").strip()
relay = (getattr(settings, "gpu_encode_relay_base_url", "") or "").strip()
relay_internal = (getattr(settings, "gpu_encode_relay_internal_base_url", "") or "").strip()
if not endpoint or not relay:
return None
return 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),
pre_warm=getattr(settings, "gpu_encode_pre_warm", True),
post_first_byte_timeout=getattr(settings, "gpu_encode_post_first_byte_timeout", 20.0),
vcodec=getattr(settings, "gpu_encode_vcodec", "h264_nvenc"),
preset=getattr(settings, "gpu_encode_preset", "p4"),
crf=getattr(settings, "gpu_encode_crf", 23),
bitrate=getattr(settings, "gpu_encode_bitrate", "") or "",
relay_secret=getattr(settings, "gpu_encode_relay_secret", "") or "",
oss_tmp_prefix=getattr(settings, "gpu_encode_oss_tmp_prefix", "tmp/gpu-mezzanine/"),
)
def get_gpu_encoder() -> Optional[GpuEncoderClient]:
"""返回进程级单例;未启用或未配置返回 None。"""
global _default_client, _default_client_initialized
if not _default_client_initialized:
_default_client_initialized = True
try:
_default_client = _build_client_from_settings()
except Exception as e: # noqa: BLE001
logger.warning("[gpu-encoder] failed to init client (CPU fallback): %s", e)
_default_client = None
return _default_client
def reset_gpu_encoder_for_tests() -> None:
global _default_client, _default_client_initialized
_default_client = None
_default_client_initialized = False
def is_gpu_encode_enabled() -> bool:
return get_gpu_encoder() is not None