e4723bfb1b
CI/CD Pipeline / Check if frontend-only change (push) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (push) Successful in 1s
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 2s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
CI/CD Pipeline / Frontend Lint (push) Has been skipped
CI/CD Pipeline / PR Build API Image (push) Has been skipped
CI/CD Pipeline / PR Build Web Image (push) Has been skipped
CI/CD Pipeline / PR Build Worker Image (push) Has been skipped
CI/CD Pipeline / Validate - Style (pull_request) Has been skipped
CI/CD Pipeline / Validate - Security (pull_request) Has been skipped
CI/CD Pipeline / Validate - Python (mypy + alembic) (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 / Check push changed paths (push) Successful in 13s
CI/CD Pipeline / Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / Unit Tests (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 / 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 / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 55s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 1m14s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Build Staging API Image (push) Successful in 1m10s
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 / Canary Release to Production (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 40s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m31s
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 2m56s
CI/CD Pipeline / Build Staging Web Image (push) Successful in 2m24s
CI/CD Pipeline / Retag skipped Staging API Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Has been skipped
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m27s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been skipped
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Failing after 33s
CI/CD Pipeline / Staging E2E Tests (push) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (push) Has been skipped
CI/CD Pipeline / ACR Image Cleanup (push) Has been skipped
CI/CD Pipeline / Validate - Style (push) Successful in 3m53s
CI/CD Pipeline / Frontend Unit Tests (push) Failing after 4m47s
CI/CD Pipeline / Integration Tests (push) Successful in 5m28s
AI Code Review / AI Code Review (pull_request) Successful in 6m51s
CI/CD Pipeline / Validate - Security (push) Successful in 7m4s
CI/CD Pipeline / PR Build Web Image (pull_request) Failing after 9m19s
CI/CD Pipeline / CI Gate (pull_request) Failing after 1s
CI/CD Pipeline / Unit Tests (push) Successful in 10m32s
CI/CD Pipeline / Build Production API Image (push) Has been skipped
CI/CD Pipeline / CI Gate (push) Has been skipped
CI/CD Pipeline / Build Production Worker Image (push) Has been skipped
CI/CD Pipeline / Build Production Web Image (push) Has been skipped
CI/CD Pipeline / Canary Release to Production (push) Has been skipped
CI/CD Pipeline / Deploy Production (push) Has been skipped
CI/CD Pipeline / Production Browser E2E (push) Has been skipped
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com> Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
216 lines
7.5 KiB
Python
216 lines
7.5 KiB
Python
"""GPU 编码回传 relay 端点。
|
||
|
||
两个用途:
|
||
1. 结果回传(原):P4000 编码完成后通过 HTTP PUT 把结果 mp4 写到 /{key};Worker 用同 URL GET 回本地。
|
||
2. Mezzanine 中转(新):Worker 先把 CPU ultrafast 编码出的 mezzanine 通过 PUT 到 /mezzanine/{key},
|
||
P4000 通过 Tailscale 内网直接 GET 下载,跳过公网 OSS 中转,节省 18-20s 固定延迟。
|
||
编码完成后 DELETE 清理。
|
||
|
||
安全:
|
||
- 生产环境必须配置 GPU_ENCODE_RELAY_SECRET;token=xxx 查询参数必须匹配。
|
||
- key 为随机 hex,无法被枚举。
|
||
- 写入/读取后 worker 会调用 DELETE 主动清理;文件落地在 generated-files/gpu_relay/。
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
import os
|
||
import secrets
|
||
import time
|
||
import uuid
|
||
from pathlib import Path
|
||
from typing import Optional
|
||
|
||
from fastapi import APIRouter, HTTPException, Query, Request
|
||
from fastapi.responses import FileResponse, Response
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
router = APIRouter(prefix="/internal/gpu-relay", tags=["Internal-GpuRelay"])
|
||
|
||
_DEFAULT_SECRET_LOGGED = False
|
||
|
||
|
||
def _relay_dir() -> Path:
|
||
base = os.getenv("GENERATED_FILES_DIR", "/app/generated")
|
||
sub = os.getenv("GPU_ENCODE_RELAY_DIR", "gpu_relay")
|
||
p = Path(base) / sub
|
||
p.mkdir(parents=True, exist_ok=True)
|
||
return p
|
||
|
||
|
||
def _mezzanine_dir() -> Path:
|
||
p = _relay_dir() / "mezzanine"
|
||
p.mkdir(parents=True, exist_ok=True)
|
||
return p
|
||
|
||
|
||
def _secret() -> str:
|
||
global _DEFAULT_SECRET_LOGGED
|
||
secret = (os.getenv("GPU_ENCODE_RELAY_SECRET", "") or "").strip()
|
||
if not secret:
|
||
env = (os.getenv("APP_ENV", os.getenv("ENV", "development"))).lower()
|
||
if env in ("production", "prod"):
|
||
raise RuntimeError("GPU_ENCODE_RELAY_SECRET must be set in production")
|
||
secret = os.environ.setdefault("GPU_ENCODE_RELAY_SECRET", secrets.token_urlsafe(32))
|
||
if not _DEFAULT_SECRET_LOGGED:
|
||
logger.warning(
|
||
"[gpu-relay] GPU_ENCODE_RELAY_SECRET not set; using ephemeral dev token (%s...)",
|
||
secret[:8],
|
||
)
|
||
_DEFAULT_SECRET_LOGGED = True
|
||
return secret
|
||
|
||
|
||
def _safe_key(key: str) -> str:
|
||
"""只允许合法文件名字符,防 path traversal。"""
|
||
k = key.strip()
|
||
if not k or "/" in k or "\\" in k or k in (".", "..") or not all(
|
||
c.isalnum() or c in "-_" for c in k
|
||
):
|
||
raise HTTPException(status_code=400, detail="invalid key")
|
||
return k
|
||
|
||
|
||
def _check_token(tok: Optional[str]) -> None:
|
||
if not tok or tok != _secret():
|
||
raise HTTPException(status_code=401, detail="unauthorized")
|
||
|
||
|
||
async def _atomic_write(request: Request, dst: Path, log_prefix: str, key_for_log: str) -> int:
|
||
"""通用原子写入(流式 → .part → replace)。返回字节数。"""
|
||
tmp = dst.with_suffix(dst.suffix + ".part")
|
||
size = 0
|
||
t0 = time.time()
|
||
try:
|
||
with open(tmp, "wb") as f:
|
||
async for chunk in request.stream():
|
||
f.write(chunk)
|
||
size += len(chunk)
|
||
os.replace(tmp, dst)
|
||
except Exception as e: # noqa: BLE001
|
||
if tmp.exists():
|
||
try:
|
||
tmp.unlink()
|
||
except OSError:
|
||
pass
|
||
logger.exception("[gpu-relay] %s PUT failed key=%s", log_prefix, key_for_log)
|
||
raise HTTPException(status_code=500, detail=f"write failed: {e}") from e
|
||
logger.info(
|
||
"[gpu-relay] %s PUT key=%s size=%d took=%.2fs",
|
||
log_prefix, key_for_log, size, time.time() - t0,
|
||
)
|
||
return size
|
||
|
||
|
||
def _file_response(path: Path, download_name: str) -> FileResponse:
|
||
if not path.exists():
|
||
raise HTTPException(status_code=404, detail="not found")
|
||
return FileResponse(path=path, media_type="video/mp4", filename=f"{download_name}.mp4")
|
||
|
||
|
||
def _head_response(path: Path) -> Response:
|
||
if not path.exists():
|
||
return Response(status_code=404)
|
||
return Response(
|
||
status_code=200,
|
||
media_type="video/mp4",
|
||
headers={"Content-Length": str(path.stat().st_size)},
|
||
)
|
||
|
||
|
||
def _safe_delete(path: Path, err_detail: str) -> dict:
|
||
try:
|
||
if path.exists():
|
||
path.unlink()
|
||
except OSError as e:
|
||
raise HTTPException(status_code=500, detail=f"{err_detail}: {e}") from e
|
||
return {"ok": True}
|
||
|
||
|
||
# ── Worker 侧 URL 构造 ─────────────────────────────────────────────────
|
||
def build_relay_put_url(base_url: str, key: str, secret: str) -> str:
|
||
"""给 P4000 回传结果用的 PUT URL(外部/Tailscale 可达)。"""
|
||
return f"{base_url.rstrip('/')}/api/v1/internal/gpu-relay/{key}?token={secret}"
|
||
|
||
|
||
def build_relay_get_url(base_url: str, key: str, secret: str) -> str:
|
||
"""Worker 取回结果用的 GET URL。"""
|
||
return build_relay_put_url(base_url, key, secret)
|
||
|
||
|
||
def build_mezzanine_put_url(base_url: str, key: str, secret: str) -> str:
|
||
"""Worker 上传 mezzanine 用的 PUT URL(Docker 内网或 Tailscale)。"""
|
||
return f"{base_url.rstrip('/')}/api/v1/internal/gpu-relay/mezzanine/{key}?token={secret}"
|
||
|
||
|
||
def build_mezzanine_get_url(base_url: str, key: str, secret: str) -> str:
|
||
"""P4000 下载 mezzanine 用的 GET URL(必须是 P4000 可达地址,通常是 Tailscale host:8092)。"""
|
||
return build_mezzanine_put_url(base_url, key, secret)
|
||
|
||
|
||
def generate_key() -> str:
|
||
return uuid.uuid4().hex
|
||
|
||
|
||
# ── 编码结果:PUT/GET/HEAD/DELETE /{key} ──────────────────────────────
|
||
@router.put("/{key}")
|
||
async def put_object(key: str, request: Request, token: Optional[str] = Query(None)):
|
||
_check_token(token)
|
||
safe = _safe_key(key)
|
||
size = await _atomic_write(request, _relay_dir() / safe, "result", safe)
|
||
return {"ok": True, "key": safe, "size": size}
|
||
|
||
|
||
@router.get("/{key}")
|
||
async def get_object(key: str, token: Optional[str] = Query(None)):
|
||
_check_token(token)
|
||
safe = _safe_key(key)
|
||
return _file_response(_relay_dir() / safe, safe)
|
||
|
||
|
||
@router.head("/{key}")
|
||
async def head_object(key: str, token: Optional[str] = Query(None)):
|
||
_check_token(token)
|
||
safe = _safe_key(key)
|
||
return _head_response(_relay_dir() / safe)
|
||
|
||
|
||
@router.delete("/{key}")
|
||
async def delete_object(key: str, token: Optional[str] = Query(None)):
|
||
_check_token(token)
|
||
safe = _safe_key(key)
|
||
return _safe_delete(_relay_dir() / safe, "delete failed")
|
||
|
||
|
||
# ── Mezzanine 中转:PUT/GET/HEAD/DELETE /mezzanine/{key} ─────────────
|
||
# Worker 上传 mezzanine 用;P4000 通过 Tailscale 直接 GET 下载。
|
||
@router.put("/mezzanine/{key}")
|
||
async def put_mezzanine(key: str, request: Request, token: Optional[str] = Query(None)):
|
||
_check_token(token)
|
||
safe = _safe_key(key)
|
||
dst = _mezzanine_dir() / f"{safe}.mp4"
|
||
size = await _atomic_write(request, dst, "mezzanine", safe)
|
||
return {"ok": True, "key": safe, "size": size}
|
||
|
||
|
||
@router.get("/mezzanine/{key}")
|
||
async def get_mezzanine(key: str, token: Optional[str] = Query(None)):
|
||
_check_token(token)
|
||
safe = _safe_key(key)
|
||
return _file_response(_mezzanine_dir() / f"{safe}.mp4", f"{safe}-mezzanine")
|
||
|
||
|
||
@router.head("/mezzanine/{key}")
|
||
async def head_mezzanine(key: str, token: Optional[str] = Query(None)):
|
||
_check_token(token)
|
||
safe = _safe_key(key)
|
||
return _head_response(_mezzanine_dir() / f"{safe}.mp4")
|
||
|
||
|
||
@router.delete("/mezzanine/{key}")
|
||
async def delete_mezzanine(key: str, token: Optional[str] = Query(None)):
|
||
_check_token(token)
|
||
safe = _safe_key(key)
|
||
return _safe_delete(_mezzanine_dir() / f"{safe}.mp4", "mezzanine delete failed")
|