diff --git a/apps/worker/video_processing/unified_render_service.py b/apps/worker/video_processing/unified_render_service.py index ec3e6bfca..a92e3bbe8 100755 --- a/apps/worker/video_processing/unified_render_service.py +++ b/apps/worker/video_processing/unified_render_service.py @@ -57,13 +57,13 @@ from video_processing.tts_engine import TtsEngine from video_processing.watermark_engine import WatermarkConfig, WatermarkEngine from packages.domain.render_layer_utils import LAYER_Z_INDEX as _IMPORTED_LAYER_Z_INDEX -from packages.shared.gpu_encoder import GpuEncodeError, get_gpu_encoder from packages.domain.render_layer_utils import clip_adjusted_duration as _clip_adjusted_duration_pure from packages.domain.render_layer_utils import clip_effective_duration as _clip_effective_duration_pure from packages.domain.render_layer_utils import clip_playback_speed as _clip_playback_speed_pure from packages.domain.render_layer_utils import estimate_total_duration as _estimate_total_duration_pure from packages.domain.render_layer_utils import resolve_layer_role as _resolve_layer_role_pure from packages.domain.tts_config import TtsConfig +from packages.shared.gpu_encoder import GpuEncodeError, get_gpu_encoder logger = logging.getLogger(__name__) @@ -2192,13 +2192,15 @@ class UnifiedRenderService: if health.ready: logger.info( "[gpu-encoder] healthy endpoint=%s gpu=%s", - client.endpoint, health.gpu_name, + client.endpoint, + health.gpu_name, ) self._gpu_health_ok = True else: logger.warning( "[gpu-encoder] not ready: %s (endpoint=%s)", - health.error, client.endpoint, + health.error, + client.endpoint, ) self._gpu_health_ok = False except Exception as e: # noqa: BLE001 @@ -2254,7 +2256,8 @@ class UnifiedRenderService: return False logger.info( "[gpu-encoder] mezzanine ready: %s (%.1fs, %d bytes), dispatching to P4000 nvenc...", - mezzanine_path.name, time.time() - t0, + mezzanine_path.name, + time.time() - t0, mezzanine_path.stat().st_size if mezzanine_path.exists() else 0, ) @@ -2269,7 +2272,8 @@ class UnifiedRenderService: ) logger.info( "[gpu-encoder] GPU nvenc encode done: %s (total %.1fs)", - output_path.name, time.time() - t0, + output_path.name, + time.time() - t0, ) return True except GpuEncodeError as e: diff --git a/apps/worker/worker_app/tasks/_startup.py b/apps/worker/worker_app/tasks/_startup.py index 51c0fa048..ab232f874 100644 --- a/apps/worker/worker_app/tasks/_startup.py +++ b/apps/worker/worker_app/tasks/_startup.py @@ -341,20 +341,26 @@ def _probe_gpu_encoder_on_ready(sender, **kwargs): """Worker 启动完成后探测 P4000 GPU NVENC 节点状态,打日志。""" try: from packages.shared.gpu_encoder import get_gpu_encoder + client = get_gpu_encoder() if client is None: - logger.info("[gpu-encoder] disabled (ENABLE_GPU_ENCODE=false or endpoint not configured), using CPU libx264") + logger.info( + "[gpu-encoder] disabled (ENABLE_GPU_ENCODE=false or endpoint not configured), using CPU libx264" + ) return health = client.check_health() if health.ready: logger.info( "[gpu-encoder] NVENC enabled: endpoint=%s gpu=%s worker=%s", - client.endpoint, health.gpu_name, health.worker, + client.endpoint, + health.gpu_name, + health.worker, ) else: logger.warning( "[gpu-encoder] configured but NOT ready: %s (endpoint=%s) — falling back to CPU", - health.error, client.endpoint, + health.error, + client.endpoint, ) except Exception as e: # noqa: BLE001 logger.warning("[gpu-encoder] startup probe error (will retry on first job, CPU fallback): %s", e) diff --git a/packages/shared/gpu_encoder.py b/packages/shared/gpu_encoder.py index 0ac1ef128..2a4dcfbe1 100644 --- a/packages/shared/gpu_encoder.py +++ b/packages/shared/gpu_encoder.py @@ -12,12 +12,12 @@ 任何环节失败抛 GpuEncodeError,调用方应 fallback 到 CPU libx264。 """ + from __future__ import annotations import json import logging import os -import secrets import socket import time import urllib.error @@ -161,7 +161,10 @@ class GpuEncoderClient: 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"), + job.get("job_id"), + job.get("ffmpeg_rc"), + job.get("size"), + job.get("duration"), ) # 5. download result from relay to output_path @@ -173,7 +176,10 @@ class GpuEncoderClient: logger.info( "[gpu-encoder] encode ok: %s → %s (%d bytes) total=%.2fs", - mezzanine_path.name, output_path.name, size, time.time() - t_total, + mezzanine_path.name, + output_path.name, + size, + time.time() - t_total, ) return {"job": job, "output_size": size, "output_path": str(output_path)} @@ -214,7 +220,8 @@ class GpuEncoderClient: req_timeout = body.get("timeout", self.sync_timeout) + 60 payload = json.dumps(body).encode("utf-8") req = urllib.request.Request( - url, data=payload, + url, + data=payload, headers={"Content-Type": "application/json"}, method="POST", ) @@ -267,8 +274,10 @@ class GpuEncoderClient: return size except (urllib.error.URLError, socket.timeout, TimeoutError, ConnectionError) as e: if tmp.exists(): - try: tmp.unlink() - except OSError: pass + try: + tmp.unlink() + except OSError: + pass raise GpuEncodeError(f"failed to download from relay: {e}") from e def _relay_delete(self, url: str) -> None: @@ -303,6 +312,7 @@ class GpuEncoderClient: 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) @@ -319,6 +329,7 @@ _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 diff --git a/tests/unit/test_gpu_encoder.py b/tests/unit/test_gpu_encoder.py index 294617758..1c72190e1 100644 --- a/tests/unit/test_gpu_encoder.py +++ b/tests/unit/test_gpu_encoder.py @@ -1,4 +1,5 @@ """GpuEncoderClient 单元测试:mock HTTP,验证 health/sync/fallback 逻辑。""" + from __future__ import annotations import json @@ -100,7 +101,12 @@ class TestPostSync: } with mock.patch("urllib.request.urlopen", return_value=_fake_response(body=result_body)) as m: res = client._post_sync( - {"inputs": {"in.mp4": "http://x"}, "ffmpeg_args": ["-i", "in.mp4"], "output_url": "http://relay/k?token=s", "timeout": 30}, + { + "inputs": {"in.mp4": "http://x"}, + "ffmpeg_args": ["-i", "in.mp4"], + "output_url": "http://relay/k?token=s", + "timeout": 30, + }, mezzanine_path=Path("/tmp/fake.mp4"), ) assert res["status"] == "completed" @@ -113,7 +119,9 @@ class TestPostSync: body = {"status": "failed", "ffmpeg_rc": 1, "message": "Invalid data found"} with mock.patch("urllib.request.urlopen", return_value=_fake_response(body=body)): with pytest.raises(GpuEncodeError, match="ffmpeg_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}, mezzanine_path=Path("/tmp/x") + ) def test_http_4xx_raises(self, client): err = urllib.error.HTTPError( @@ -121,7 +129,9 @@ 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}, mezzanine_path=Path("/tmp/x") + ) class TestRelayUrl: @@ -169,4 +179,5 @@ class TestDownloadToFile: class TestFfmpegOutputToMezzanineIntegration: """_ffmpeg_output_to_mezzanine is on UnifiedRenderService; unit-tested there via mocks.""" + pass