114ce95c04
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 1s
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
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 / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 23s
CI/CD Pipeline / PR Build Web Image (pull_request) Has been skipped
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 1m35s
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
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 Worker Image (pull_request) Successful in 1m55s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m12s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 3m30s
AI Code Review / AI Code Review (pull_request) Successful in 6m51s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 9m39s
CI/CD Pipeline / Validate - Style (pull_request) Failing after 10m37s
CI/CD Pipeline / Build Production API Image (pull_request) Has been cancelled
CI/CD Pipeline / Build Production Web Image (pull_request) Has been cancelled
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been cancelled
CI/CD Pipeline / Deploy Production (pull_request) Has been cancelled
CI/CD Pipeline / Production Browser E2E (pull_request) Has been cancelled
CI/CD Pipeline / Canary Release to Production (pull_request) Has been cancelled
CI/CD Pipeline / CI Gate (pull_request) Has been cancelled
CI/CD Pipeline / Validate - Security (pull_request) Has been cancelled
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Has been cancelled
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been cancelled
P0 issue: 2026-09-27 staging job a15a4667 observed 149.6s delay between worker POST and P4000 starting mezzanine download (nginx access log confirmed GET mezz at 14:12:55 vs POST at 14:10:26). Root cause is server/network-side (likely Tailscale DERP hole-punching or httpx connection pool rebuild after long idle), hard to reproduce on demand. This client-side fix adds two layers of defense: 1. pre_warm: before POST, if idle >60s since last successful comm, send a GET /health with 3s timeout to warm the Tailscale path. Failure is non-fatal; POST proceeds anyway. 2. post_first_byte_timeout (default 20s): split POST timeout into two phases — first byte uses short timeout so a stalled link fails fast and caller can CPU-fallback; after first byte, relax to ffmpeg_timeout+60s for real encode+upload. Both configurable via settings.gpu_encode_pre_warm / settings.gpu_encode_post_first_byte_timeout. Also updates check_health() to refresh _last_ok_ts on healthy response.
565 lines
26 KiB
Python
565 lines
26 KiB
Python
"""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,
|
||
)
|
||
# 剩余用长超时
|
||
resp.fp._sock.settimeout(ffmpeg_timeout + 60)
|
||
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"
|
||
raise GpuEncodeError(
|
||
f"P4000 {hint} timeout 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
|