feat(gpu): 接入 P4000 NVENC 硬件编码加速
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 1s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
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 Web Image (pull_request) 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 / 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 API Image (pull_request) Successful in 51s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 1m22s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 1m50s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 1m59s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m0s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 4m56s
AI Code Review / AI Code Review (pull_request) Successful in 8m14s
CI/CD Pipeline / PR Build Worker Image (pull_request) Failing after 9m8s
CI/CD Pipeline / Validate - Style (pull_request) Failing after 9m43s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 11m19s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 40m55s
CI/CD Pipeline / Build Production 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 / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / CI Gate (pull_request) Failing after 2s
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 1s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
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 Web Image (pull_request) 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 / 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 API Image (pull_request) Successful in 51s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 1m22s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 1m50s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 1m59s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m0s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 4m56s
AI Code Review / AI Code Review (pull_request) Successful in 8m14s
CI/CD Pipeline / PR Build Worker Image (pull_request) Failing after 9m8s
CI/CD Pipeline / Validate - Style (pull_request) Failing after 9m43s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 11m19s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 40m55s
CI/CD Pipeline / Build Production 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 / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / CI Gate (pull_request) Failing after 2s
- packages/config/base.py: 新增 GPU_ENCODE_* 配置项(开关/endpoint/relay/secret/编码参数/fallback)
- packages/shared/gpu_encoder.py: 新增 GpuEncoderClient,封装 health 探测 + mezzanine 上传 OSS + P4000 nvenc 编码 + relay 回传下载,失败抛 GpuEncodeError 触发 CPU 降级
- apps/api/app/api/routes/gpu_relay.py: 新增 internal PUT/GET/DELETE /api/v1/internal/gpu-relay/{key}(token 鉴权),P4000 PUT 编码结果,worker GET 下载
- apps/api/app/api/router.py: 注册 gpu_relay_router
- apps/worker/video_processing/unified_render_service.py: _execute_ffmpeg 和 _render_pass_through 尝试 GPU 路径:CPU ultrafast mezzanine → P4000 nvenc → 输出到最终路径;任何失败自动回退到原 CPU libx264 路径
- apps/worker/worker_app/tasks/_startup.py: worker_ready 时探测 P4000 健康并打日志
- tests/unit/test_gpu_encoder.py: GpuEncoderClient 单测(health/sync 调用/失败/fallback/relay URL)
架构:
- 输入:CPU 输出 libx264 ultrafast mezzanine → 上传 OSS 临时前缀 → P4000 签名 URL 下载
- 输出:P4000 PUT → 宿主机 nginx(tailscale:80)→ API /api/ 反代 → gpu_relay 路由落盘到 generated/gpu_relay/
- 回传:worker 通过 docker 网络 http://xiaoxia-api-staging:8000 GET 下载最终 mp4 到 output_path
- 降级:GPU 任何环节异常(health/上传/编码/回传/下载)→ 原 CPU 路径继续执行,不影响成片
This commit is contained in:
@@ -0,0 +1,172 @@
|
||||
"""GpuEncoderClient 单元测试:mock HTTP,验证 health/sync/fallback 逻辑。"""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import tempfile
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
from http.client import HTTPResponse
|
||||
from io import BytesIO
|
||||
from pathlib import Path
|
||||
from unittest import mock
|
||||
|
||||
import pytest
|
||||
|
||||
from packages.shared.gpu_encoder import (
|
||||
GpuEncodeError,
|
||||
GpuEncoderClient,
|
||||
GpuHealth,
|
||||
reset_gpu_encoder_for_tests,
|
||||
)
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _reset_singleton():
|
||||
reset_gpu_encoder_for_tests()
|
||||
yield
|
||||
reset_gpu_encoder_for_tests()
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def client():
|
||||
return GpuEncoderClient(
|
||||
endpoint="http://gpu.example.com:8900",
|
||||
relay_base_url="http://api.example.com",
|
||||
sync_timeout=60,
|
||||
health_timeout=2,
|
||||
relay_secret="test-secret",
|
||||
)
|
||||
|
||||
|
||||
def _fake_response(status: int = 200, body: dict | bytes | None = None, headers=None):
|
||||
if isinstance(body, dict):
|
||||
data = json.dumps(body).encode("utf-8")
|
||||
elif body is None:
|
||||
data = b""
|
||||
else:
|
||||
data = body
|
||||
resp = mock.MagicMock(spec=HTTPResponse)
|
||||
resp.status = status
|
||||
resp.read.return_value = data
|
||||
resp.__enter__ = mock.MagicMock(return_value=resp)
|
||||
resp.__exit__ = mock.MagicMock(return_value=False)
|
||||
return resp
|
||||
|
||||
|
||||
class TestHealthCheck:
|
||||
def test_healthy_nvenc_available(self, client):
|
||||
body = {
|
||||
"status": "healthy",
|
||||
"worker": "gpu-worker-1",
|
||||
"gpu": {"name": "Quadro P4000"},
|
||||
"nvenc": {"h264_nvenc": True, "hevc_nvenc": True},
|
||||
}
|
||||
with mock.patch("urllib.request.urlopen", return_value=_fake_response(body=body)):
|
||||
h = client.check_health()
|
||||
assert h.healthy
|
||||
assert h.nvenc_h264
|
||||
assert h.ready
|
||||
assert h.gpu_name == "Quadro P4000"
|
||||
|
||||
def test_connection_error_returns_unhealthy(self, client):
|
||||
with mock.patch("urllib.request.urlopen", side_effect=urllib.error.URLError("timeout")):
|
||||
h = client.check_health()
|
||||
assert not h.healthy
|
||||
assert "health probe failed" in h.error
|
||||
|
||||
def test_bad_json_returns_unhealthy(self, client):
|
||||
with mock.patch("urllib.request.urlopen", return_value=_fake_response(body=b"not json")):
|
||||
h = client.check_health()
|
||||
assert not h.healthy
|
||||
|
||||
def test_nvenc_unavailable(self, client):
|
||||
body = {"status": "healthy", "gpu": {"name": "test"}, "nvenc": {"h264_nvenc": False}}
|
||||
with mock.patch("urllib.request.urlopen", return_value=_fake_response(body=body)):
|
||||
h = client.check_health()
|
||||
assert h.healthy
|
||||
assert not h.ready
|
||||
|
||||
|
||||
class TestPostSync:
|
||||
def test_completed_job_returns_dict(self, client):
|
||||
result_body = {
|
||||
"job_id": "j1",
|
||||
"status": "completed",
|
||||
"ffmpeg_rc": 0,
|
||||
"uploaded": True,
|
||||
"duration": 5.1,
|
||||
"size": 123456,
|
||||
}
|
||||
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},
|
||||
mezzanine_path=Path("/tmp/fake.mp4"),
|
||||
)
|
||||
assert res["status"] == "completed"
|
||||
assert res["ffmpeg_rc"] == 0
|
||||
# verify request sent to sync endpoint
|
||||
req = m.call_args[0][0]
|
||||
assert req.full_url == "http://gpu.example.com:8900/api/render/sync"
|
||||
|
||||
def test_ffmpeg_failure_raises(self, client):
|
||||
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"))
|
||||
|
||||
def test_http_4xx_raises(self, client):
|
||||
err = urllib.error.HTTPError(
|
||||
url="http://gpu/render/sync", code=422, msg="Unprocessable", hdrs={}, fp=BytesIO(b"bad request")
|
||||
)
|
||||
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"))
|
||||
|
||||
|
||||
class TestRelayUrl:
|
||||
def test_url_contains_token_and_key(self, client):
|
||||
url = client._relay_url("abc123", "secret!")
|
||||
assert "abc123" in url
|
||||
assert "token=secret%21" in url # urlencoded
|
||||
assert url.startswith("http://api.example.com/api/v1/internal/gpu-relay/")
|
||||
|
||||
|
||||
class TestGetRelaySecret:
|
||||
def test_explicit_secret_used(self, client):
|
||||
assert client._get_relay_secret() == "test-secret"
|
||||
|
||||
def test_env_secret_used_when_not_explicit(self, monkeypatch):
|
||||
monkeypatch.setenv("GPU_ENCODE_RELAY_SECRET", "from-env")
|
||||
monkeypatch.setenv("APP_ENV", "staging")
|
||||
c = GpuEncoderClient(endpoint="http://gpu", relay_base_url="http://api")
|
||||
assert c._get_relay_secret() == "from-env"
|
||||
|
||||
def test_prod_without_secret_raises(self, monkeypatch):
|
||||
monkeypatch.delenv("GPU_ENCODE_RELAY_SECRET", raising=False)
|
||||
monkeypatch.setenv("APP_ENV", "production")
|
||||
c = GpuEncoderClient(endpoint="http://gpu", relay_base_url="http://api")
|
||||
with pytest.raises(GpuEncodeError, match="GPU_ENCODE_RELAY_SECRET"):
|
||||
c._get_relay_secret()
|
||||
|
||||
|
||||
class TestDownloadToFile:
|
||||
def test_writes_file(self, client, tmp_path):
|
||||
data = b"hello" * 1000
|
||||
out = tmp_path / "out.mp4"
|
||||
with mock.patch("urllib.request.urlopen", return_value=_fake_response(body=data)):
|
||||
size = client._download_to_file("http://relay/k?token=s", out)
|
||||
assert size == len(data)
|
||||
assert out.read_bytes() == data
|
||||
|
||||
def test_empty_file_raises(self, client, tmp_path):
|
||||
out = tmp_path / "out.mp4"
|
||||
with mock.patch("urllib.request.urlopen", return_value=_fake_response(body=b"")):
|
||||
with pytest.raises(GpuEncodeError, match="empty file"):
|
||||
client._download_to_file("http://relay/k", out)
|
||||
assert not out.exists()
|
||||
|
||||
|
||||
class TestFfmpegOutputToMezzanineIntegration:
|
||||
"""_ffmpeg_output_to_mezzanine is on UnifiedRenderService; unit-tested there via mocks."""
|
||||
pass
|
||||
Reference in New Issue
Block a user