"""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. 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 从 relay GET mezzanine → h264_nvenc 编码 → PUT 最终 mp4 到 put_url 6. 本客户端通过 get_url(Docker 内网)下载最终文件到 output_path,然后 DELETE 清理 7. 删除 relay 上的 mezzanine 临时文件(以及 OSS fallback 的 key) 任何环节失败抛 GpuEncodeError,调用方应 fallback 到 CPU libx264。 """ 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, 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.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/" 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: return 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")), ) 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. call P4000 sync render body = { "inputs": {"in.mp4": input_url}, "ffmpeg_args": ffmpeg_args, "output_url": put_url, "timeout": int(timeout), } job = self._post_sync(body) 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 _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") req = urllib.request.Request( url, data=payload, headers={"Content-Type": "application/json"}, method="POST", ) t0 = time.time() try: with urllib.request.urlopen(req, timeout=req_timeout) as resp: raw = resp.read().decode("utf-8") 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) as e: raise GpuEncodeError(f"P4000 connection error: {e}") from e 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), 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