"""P4000 NVENC 远程编码客户端。 完整链路(encode_video_file): 1. CPU 滤镜已在本地生成 mezzanine 中间片(libx264 ultrafast) 2. 上传 mezzanine 到 OSS 临时前缀,拿到签名 GET URL 3. 生成 relay 一次性 key,构造两个带 token 的 URL: - put_url:给 P4000 回传结果,走 relay_base_url(外部可达,通常是 host:port 经 nginx) - 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 编码完成后 PUT 最终 mp4 到 put_url,API 服务落盘到 /app/generated/gpu_relay/ 6. 本客户端通过 get_url(Docker 内网)下载最终文件到 output_path,然后 DELETE 清理 7. 删除 OSS 临时 mezzanine 任何环节失败抛 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 = "", 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.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, key: str, secret: str) -> str: return f"{base_url}{self.RELAY_PATH_PREFIX}/{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_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) # ------------------------------------------------------------------ # 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。 extra_video_args: -i 之后、-c:v 之前插入的 ffmpeg 参数(如分辨率/帧率调整)。 audio_args: 音频编码参数(如 ["-c:a","aac","-b:a","128k"]);None 表示 -an 无音频。 """ 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 relay_key: Optional[str] = None 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 # 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, mezzanine_path=mezzanine_path) logger.info( "[gpu-encoder] P4000 done: job_id=%s rc=%s size=%s dur=%ss", job.get("job_id"), job.get("ffmpeg_rc"), job.get("size"), job.get("duration"), ) # 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) logger.info( "[gpu-encoder] encode ok: %s → %s (%d bytes) total=%.2fs", mezzanine_path.name, output_path.name, size, time.time() - t_total, ) return {"job": job, "output_size": size, "output_path": str(output_path)} except GpuEncodeError: raise except Exception as e: # noqa: BLE001 raise GpuEncodeError(f"unexpected: {e}") from e finally: # 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 # ------------------------------------------------------------------ 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]: 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}") # 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 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) # ------------------------------------------------------------------ # OSS helpers (optional - storage may not be available in all envs) # ------------------------------------------------------------------ def _upload_mezzanine(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 # Generate signed GET URL (1h expiry) 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, 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 # Convenience def is_gpu_encode_enabled() -> bool: return get_gpu_encoder() is not None