Compare commits

...

6 Commits

Author SHA1 Message Date
业务服务器运维 ba8f786b60 fix: 修复 _storage_key 无法写入 EditPlanClip(slots=True)
EditPlanClip 使用 @dataclass(slots=True) 不允许动态添加属性,
改为将 _storage_key 存入 clip.config 字典,在 _resolve_clips
中传递到 ResolvedClip.config,gpu_direct_pipeline 从 config 读取。
2026-09-28 14:43:17 +08:00
业务服务器运维 b04a803655 feat(worker): 全GPU直连渲染,取消mezzanine CPU中间编码
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 1s
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 1s
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
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 2m1s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m19s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 2m55s
AI Code Review / AI Code Review (pull_request) Successful in 6m56s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 6m51s
CI/CD Pipeline / Build Production API Image (pull_request) Has been cancelled
CI/CD Pipeline / Build Production Web Image (pull_request) Has been cancelled
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been cancelled
CI/CD Pipeline / Deploy Production (pull_request) Has been cancelled
CI/CD Pipeline / Production Browser E2E (pull_request) Has been cancelled
CI/CD Pipeline / Canary Release to Production (pull_request) Has been cancelled
CI/CD Pipeline / CI Gate (pull_request) Has been cancelled
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Has been cancelled
CI/CD Pipeline / Validate - Style (pull_request) Has been cancelled
CI/CD Pipeline / Validate - Security (pull_request) Has been cancelled
CI/CD Pipeline / Unit Tests (pull_request) Has been cancelled
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been cancelled
- gpu_encoder 新增 render_inputs_to_output:多输入+filter_complex 直接交 P4000 一次出片
- 新增 gpu_direct_pipeline:原始素材/BGM 签名URL、TTS 上传、drawtext 标题/字幕、
  filter_complex 完成 trim/scale/pad/concat/边缘crop/amix,末端 h264_nvenc 仅编码一次
- unified_render_service:render() 步骤4.8 接入直连,命中即出片;
  不满足条件或 GPU 失败自动回退现有 mezzanine/CPU 链路;config gpu_direct_enabled 可灰度关闭
- render_adapter:_download_assets 额外返回 asset_storage_map,
  _do_render 给 clip 注入 _storage_key 供直连签名

消除 mezzanine 85s + 边缘裁剪 26s + 合成 2s + 中转 1s ≈ 114s CPU 开销
2026-09-28 13:30:31 +08:00
xiaoxia e496f127a3 feat(p1): OSS 双endpoint分离 — 上传/下载走VPC内网,签名URL走公网 (#2085)
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 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 2s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 3s
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 / 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 / Validate - Security (pull_request) Has been skipped
CI/CD Pipeline / Validate - Style (pull_request) Has been skipped
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Has been skipped
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 / PR Build API Image (pull_request) Successful in 1m25s
CI/CD Pipeline / Build Staging API Image (push) Successful in 1m24s
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 / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (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 / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (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
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m2s
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 2m58s
CI/CD Pipeline / PR Build Web Image (pull_request) Successful in 2m55s
CI/CD Pipeline / Build Staging Web Image (push) Successful in 3m6s
CI/CD Pipeline / Integration Tests (push) Successful in 3m33s
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 2m20s
CI/CD Pipeline / Validate - Style (push) Successful in 4m8s
CI/CD Pipeline / Retag skipped Staging Web Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging API Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (push) Failing after 5m0s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 1m4s
CI/CD Pipeline / Validate - Security (push) Successful in 6m31s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 6m53s
CI/CD Pipeline / ACR Image Cleanup (push) Successful in 1m43s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been skipped
AI Code Review / AI Code Review (pull_request) Successful in 7m2s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 2m56s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 3m44s
CI/CD Pipeline / PR Build Worker Image (pull_request) Failing after 9m36s
CI/CD Pipeline / CI Gate (pull_request) Failing after 1s
CI/CD Pipeline / Unit Tests (push) Successful in 10m39s
CI/CD Pipeline / Build Production API Image (push) Has been skipped
CI/CD Pipeline / Build Production Web Image (push) Has been skipped
CI/CD Pipeline / Build Production Worker Image (push) Has been skipped
CI/CD Pipeline / CI Gate (push) Has been skipped
CI/CD Pipeline / Deploy Production (push) Has been skipped
CI/CD Pipeline / Production Browser E2E (push) Has been skipped
CI/CD Pipeline / Canary Release to Production (push) Has been skipped
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com>
Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
2026-09-28 11:35:13 +08:00
xiaoxia 423be1446f Merge pull request 'fix(p0): GPU mezzanine传输切回OSS绕开Tailscale反向卡顿' (#2084) from fix/p0-gpu-mezzanine-oss into develop
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 / Frontend Lint (push) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 3s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 3s
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 / PR Build API Image (push) 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 / 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 / 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 / Check push changed paths (push) Successful in 13s
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 API Image (pull_request) Successful in 1m34s
CI/CD Pipeline / Build Staging API Image (push) Successful in 1m41s
CI/CD Pipeline / Integration Tests (push) Successful in 2m3s
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
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m11s
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / ACR Image Cleanup (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 3m10s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 3m2s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m17s
CI/CD Pipeline / PR Build Web Image (pull_request) Successful in 3m9s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been skipped
CI/CD Pipeline / CI Gate (pull_request) Successful in 3s
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 1m18s
CI/CD Pipeline / Build Staging Web Image (push) Successful in 1m43s
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
CI/CD Pipeline / Frontend Unit Tests (push) Failing after 4m36s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 1m11s
CI/CD Pipeline / Validate - Style (push) Successful in 5m30s
CI/CD Pipeline / Validate - Security (push) Successful in 5m41s
CI/CD Pipeline / ACR Image Cleanup (push) Successful in 1m41s
AI Code Review / AI Code Review (pull_request) Successful in 6m46s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 3m51s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 5m8s
CI/CD Pipeline / Unit Tests (push) Successful in 10m44s
CI/CD Pipeline / Build Production API Image (push) Has been skipped
CI/CD Pipeline / Build Production Web Image (push) Has been skipped
CI/CD Pipeline / Build Production Worker Image (push) Has been skipped
CI/CD Pipeline / CI Gate (push) Has been skipped
CI/CD Pipeline / Deploy Production (push) Has been skipped
CI/CD Pipeline / Canary Release to Production (push) Has been skipped
CI/CD Pipeline / Production Browser E2E (push) Has been skipped
fix(p0): GPU mezzanine传输切回OSS绕开Tailscale反向卡顿

Tailscale 反向链路(P4000 从 116 relay 下载 mezzanine)间歇性 TCP stall,
导致 GPU 编码首字节 20s 超时 fallback 到 CPU 软编,渲染变慢到 ~3 分钟。

改 GPU_ENCODE_MEZZANINE_TRANSPORT=relay → oss:P4000 直接从阿里云 OSS
公网下载 mezzanine,绕开不稳定的 Tailscale 反向链路;编码结果回传仍走
relay(116→P4000 正向 POST 正常)。仅改 .env.staging,部署后 watchtower
自动拉起容器生效。
2026-09-28 08:51:37 +08:00
xiaoxia 6d9d2e8179 fix(p0): GPU mezzanine传输切回OSS绕开Tailscale反向卡顿
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 52s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 1m28s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 2m3s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 2m17s
CI/CD Pipeline / Validate - Style (pull_request) Successful in 2m52s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m6s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 4m28s
AI Code Review / AI Code Review (pull_request) Successful in 6m37s
CI/CD Pipeline / PR Build Worker Image (pull_request) Failing after 7m34s
CI/CD Pipeline / Unit Tests (pull_request) Successful in 9m11s
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 / CI Gate (pull_request) Failing after 1s
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
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 6m17s
Preview Cleanup / Cleanup Preview Environment (pull_request) Successful in 27s
ACR Cleanup / ACR Image Cleanup (pull_request_target) Successful in 4m31s
Tailscale 反向链路(P4000 从 116 relay 下载 mezzanine)间歇性 TCP stall,
导致 GPU 编码 20s 首字节超时 fallback 到 CPU 软编,渲染变慢到 ~3 分钟。

改 transport=oss:P4000 直接从阿里云 OSS 公网下载 mezzanine,绕开不稳定的
Tailscale 反向链路;编码结果回传仍走 relay(116→P4000 正向链路正常)。
2026-09-28 08:31:20 +08:00
xiaoxia b233529eee Merge pull request 'fix(ci): #2081 follow-up——CI 部署脚本修复(compose.yml 由 CI scp + --env-file)' (#2082) from fix/ci-worker-compose-followup into develop
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (push) Successful in 1s
CI/CD Pipeline / Check if frontend-only change (push) Has been skipped
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 3s
CI/CD Pipeline / Frontend Lint (push) Has been skipped
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 5s
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 / 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 / Validate - Style (pull_request) Has been skipped
CI/CD Pipeline / Validate - Security (pull_request) Has been skipped
CI/CD Pipeline / Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Has been skipped
CI/CD Pipeline / Integration 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 / PR Build API Image (pull_request) Successful in 58s
CI/CD Pipeline / Build Staging API Image (push) Successful in 1m35s
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 / 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 / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / ACR Image Cleanup (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
CI/CD Pipeline / PR Build Web Image (pull_request) Successful in 2m24s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 2m27s
CI/CD Pipeline / Build Staging Web Image (push) Successful in 1m51s
CI/CD Pipeline / CI Gate (pull_request) Successful in 2s
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 1m39s
CI/CD Pipeline / Retag skipped Staging API Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (push) Has been skipped
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m19s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been skipped
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 1m16s
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 4m40s
CI/CD Pipeline / Validate - Style (push) Successful in 5m13s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 5m22s
CI/CD Pipeline / Frontend Unit Tests (push) Failing after 5m19s
CI/CD Pipeline / ACR Image Cleanup (push) Successful in 2m7s
CI/CD Pipeline / Integration Tests (push) Successful in 6m41s
AI Code Review / AI Code Review (pull_request) Successful in 6m56s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 3m29s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 4m5s
CI/CD Pipeline / Validate - Security (push) Successful in 8m22s
CI/CD Pipeline / Unit Tests (push) Successful in 11m20s
CI/CD Pipeline / Build Production API Image (push) Has been skipped
CI/CD Pipeline / Build Production Web Image (push) Has been skipped
CI/CD Pipeline / Build Production Worker Image (push) Has been skipped
CI/CD Pipeline / CI Gate (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
fix(ci): worker 部署收敛到 compose——CI scp compose.yml + --env-file,消除并发/HC 配置漂移

修复 #2081 首次部署失败的两个问题:
- 去掉 curl 私有仓库(无 token 404),改由 CI workflow 在运行脚本前 scp infra/docker/compose.yml 到服务器
- docker compose 加 --env-file 显式指向 .env,确保 GENERATED_FILES_HOST_DIR 等变量加载
- 封装 compose() 函数统一调用;回滚走同样路径,compose 失败才 fallback docker run
2026-09-28 02:41:32 +08:00
9 changed files with 1082 additions and 431 deletions
@@ -0,0 +1,300 @@
"""全 GPU 直连渲染管线(P1)。
背景:旧链路 worker 先用 CPU libx264 把 filter_complex 输出成 mezzanine(1080p 约 85s),
上传后再由 P4000 NVENC 编码,渲染后还要单独跑一次随机边缘裁剪重编码(约 26s)。
本管线取消 mezzanine:把原始素材签名 URL 作为多输入直接交给 P4000,filter_complex 内
一步完成 trim/scale/pad/concat/边缘裁剪/drawtext 字幕,末端 h264_nvenc 只编码一次;
TTS/BGM 音频也在同一命令里 amix 合成。
约束(P1):
- 仅覆盖智能剪辑主流场景:单一主视频轨、全硬切、无 PiP/overlay/水印/贴纸/片头片尾/绿幕。
不满足条件时调用方回退到现有 mezzanine/CPU 链路(功能不回归)。
- 字幕先用 drawtext(P4000 装好中文字体后可再切 subtitles 滤镜烧 ASS)。
"""
from __future__ import annotations
import logging
import uuid
from pathlib import Path
from typing import Any, Optional
logger = logging.getLogger(__name__)
# drawtext 默认字体名(fontconfig 解析);P4000 装好中文字体后可在 config 指定
DEFAULT_DRAWTEXT_FONT = "sans"
def escape_drawtext_text(text: str) -> str:
"""转义 drawtext text= 中的特殊字符(ffmpeg 过滤器语法)。"""
if not text:
return ""
# 顺序重要:先转义反斜杠本身
s = text.replace("\\", "\\\\")
s = s.replace(":", "\\:")
s = s.replace("'", "\\'")
s = s.replace("%", "\\%")
s = s.replace(",", "\\,")
s = s.replace("[", "\\[").replace("]", "\\]")
s = s.replace(";", "\\;")
# 换行保留(drawtext 支持 %{...};真实换行需转成字面)
s = s.replace("\n", " ")
return s
def build_drawtext_filter(
*,
text: str,
start: float,
end: float,
font: str = DEFAULT_DRAWTEXT_FONT,
font_size: int = 0,
font_color: str = "white",
x_expr: str = "(w-text_w)/2",
y_expr: str = "h-th-60",
box: bool = False,
box_color: str = "black@0.5",
borderw: int = 0,
border_color: str = "black",
enable: bool = True,
) -> str:
"""构造单个 drawtext 滤镜字符串(不含输入/输出标签)。
[start, end] 秒的显示窗口通过 enable='between(t,...)' 控制。
"""
txt = escape_drawtext_text(text)
parts = [f"font={font}", f"text='{txt}'"]
if font_size and font_size > 0:
parts.append(f"fontsize={int(font_size)}")
parts.append(f"fontcolor={font_color}")
if box:
parts.append("box=1")
parts.append(f"boxcolor={box_color}")
if borderw and borderw > 0:
parts.append(f"borderw={int(borderw)}")
parts.append(f"bordercolor={border_color}")
parts.append(f"x={x_expr}")
parts.append(f"y={y_expr}")
if enable:
parts.append(f"enable='between(t,{start:.3f},{end:.3f})'")
return "drawtext=" + ":".join(parts)
def upload_local_audio_and_sign(
local_audio: Path,
*,
tmp_prefix: str = "tmp/gpu-direct-audio/",
expires: int = 3600,
) -> tuple[str, str]:
"""把本地音频(TTS/BGM)上传 OSS tmp 目录并签公网 GET URL。
Returns:
(signed_get_url, oss_key)
"""
from video_processing.oss_helpers import _storage # 类型: ignore
storage = _storage()
key = f"{tmp_prefix.rstrip('/')}/{uuid.uuid4().hex}{local_audio.suffix or '.mp3'}"
content_type = "audio/mpeg" if local_audio.suffix.lower() in (".mp3", ".mpeg") else "audio/mp4"
storage.upload_file(local_audio, key, content_type=content_type)
url = storage.get_download_url(key, expires)
return url, key
def sign_asset_url(storage_key: str, *, expires: int = 3600) -> str:
"""给原始素材 storage_key 签公网 GET URL(供 P4000 直接下载)。"""
from video_processing.oss_helpers import _storage # 类型: ignore
storage = _storage()
return storage.get_download_url(storage_key, expires)
# ── 直连渲染编排器 ────────────────────────────────────────────────────────
class DirectRenderPlan:
"""一次直连渲染的产物:inputs(裸文件名→URL)与完整 ffmpeg_args。"""
def __init__(self, inputs: dict[str, str], ffmpeg_args: list[str], oss_keys: list[str]):
self.inputs = inputs
self.ffmpeg_args = ffmpeg_args
self.oss_keys = oss_keys # 本次上传的临时音频 key(供事后清理)
def build_direct_render(
*,
resolved_clips: list[Any],
output_width: int,
output_height: int,
output_fps: int,
tts_audio: Optional[Path] = None,
bgm_audio: Optional[Path] = None,
title_text: str = "",
subtitle_segments: Optional[list[Any]] = None,
font: str = DEFAULT_DRAWTEXT_FONT,
vcodec: str = "h264_nvenc",
preset: str = "p4",
video_bitrate: str = "",
cq: int = 23,
edge_crop_pct: float = 0.0,
total_duration: float = 0.0,
) -> DirectRenderPlan:
"""构造 P4000 直连渲染所需的 inputs 与 ffmpeg_args。
视频:每段 trim/setpts/scale/pad/fps → concat(全硬切)→ 可选边缘 crop+scale → drawtext。
音频:concat 时丢弃原生音轨(只映射 [vfinal]),TTS/BGM 上传签名后 amix 混音。
"""
if not resolved_clips:
raise ValueError("build_direct_render: no resolved clips")
inputs: dict[str, str] = {}
oss_keys: list[str] = []
input_args: list[str] = []
fc: list[str] = [] # filter_complex 各段
n = len(resolved_clips)
# 1. 视频输入(原始素材签名 URL)
for i, clip in enumerate(resolved_clips):
sk = (getattr(clip, "config", None) or {}).get("_storage_key")
if not sk:
raise ValueError(f"clip {getattr(clip,'clip_id',i)} missing _storage_key")
fname = f"v{i}.mp4"
inputs[fname] = sign_asset_url(sk)
input_args.extend(["-i", fname])
# 2. 每个视频段预处理
pre_labels: list[str] = []
for i, clip in enumerate(resolved_clips):
filters: list[str] = []
start = float(getattr(clip, "start_time", 0) or 0)
eff = float(getattr(clip, "duration", 0) or 0)
if eff <= 0:
eff = float(getattr(clip, "actual_duration", 0) or 0)
if eff > 0:
if start > 0:
filters.append(f"trim=start={start:.3f}:duration={eff:.3f}")
else:
filters.append(f"trim=duration={eff:.3f}")
filters.append("setpts=PTS-STARTPTS")
speed = float(getattr(clip, "playback_speed", 1.0) or 1.0)
if abs(speed - 1.0) >= 1e-6:
filters.append(f"setpts=PTS/{speed:.4f}")
filters.append(f"scale={output_width}:{output_height}:force_original_aspect_ratio=decrease")
filters.append(f"pad={output_width}:{output_height}:trunc((ow-iw)/2):trunc((oh-ih)/2):black")
filters.append("setpts=PTS-STARTPTS")
filters.append(f"fps={output_fps}")
label = f"vc{i}"
fc.append(f"[{i}:v]{','.join(filters)}[{label}]")
pre_labels.append(label)
# 3. concat(全硬切;原生音频丢弃,v=1:a=0)
concat_in = "".join(f"[{l}]" for l in pre_labels)
fc.append(f"{concat_in}concat=n={n}:v=1:a=0[vcat]")
cur = "vcat"
# 4. 边缘裁剪降重(合并进同一条,不再单独重编码)
if edge_crop_pct and edge_crop_pct > 0:
p = float(edge_crop_pct)
keep = 1.0 - 2.0 * p
cw_expr = f"trunc(iw*{keep:.4f}/2)*2"
ch_expr = f"trunc(ih*{keep:.4f}/2)*2"
fc.append(
f"[{cur}]crop=w='{cw_expr}':h='{ch_expr}':x='(iw-{cw_expr})/2':y='(ih-{ch_expr})/2',"
f"scale={output_width}:{output_height}[vcrop]"
)
cur = "vcrop"
# 5. drawtext 字幕(标题整段 + ASR 逐句)
draw_filters: list[str] = []
if title_text.strip():
title_size = max(int(output_height * 0.05), 24)
draw_filters.append(
build_drawtext_filter(
text=title_text,
start=0.0,
end=max(total_duration, 0.1),
font=font,
font_size=title_size,
y_expr="h-th-40",
box=True,
)
)
sub_size = max(int(output_height * 0.045), 20)
for seg in subtitle_segments or []:
txt = getattr(seg, "text", "") or ""
if not txt.strip():
continue
draw_filters.append(
build_drawtext_filter(
text=txt,
start=float(getattr(seg, "start", 0)),
end=float(getattr(seg, "end", 0)),
font=font,
font_size=sub_size,
y_expr="h-th-60",
borderw=2,
)
)
if draw_filters:
prev = cur
for idx, df in enumerate(draw_filters):
out_l = "vfinal" if idx == len(draw_filters) - 1 else f"vd{idx}"
fc.append(f"[{prev}]{df}[{out_l}]")
prev = out_l
vfinal_label = prev
else:
fc.append(f"[{cur}]format=yuv420p[vfinal]")
vfinal_label = "vfinal"
# 6. 音频输入与混音
audio_items: list[tuple[int, float]] = [] # (input_index, volume)
next_idx = n
if tts_audio and Path(tts_audio).exists():
turl, tkey = upload_local_audio_and_sign(Path(tts_audio))
tname = "tts" + (Path(tts_audio).suffix or ".mp3")
inputs[tname] = turl
oss_keys.append(tkey)
input_args.extend(["-i", tname])
audio_items.append((next_idx, 1.0))
next_idx += 1
if bgm_audio and Path(bgm_audio).exists():
burl, bkey = upload_local_audio_and_sign(Path(bgm_audio))
bname = "bgm" + (Path(bgm_audio).suffix or ".mp3")
inputs[bname] = burl
oss_keys.append(bkey)
input_args.extend(["-i", bname])
audio_items.append((next_idx, 0.35))
next_idx += 1
maps: list[str] = ["-map", f"[{vfinal_label}]"]
if audio_items:
mix_labels: list[str] = []
for k, (idx, vol) in enumerate(audio_items):
alabel = f"au{k}"
fc.append(
f"[{idx}:a]aresample=44100,volume={vol:.2f},"
f"aformat=sample_fmts=fltp:channel_layouts=stereo[{alabel}]"
)
mix_labels.append(alabel)
mix_in = "".join(f"[{l}]" for l in mix_labels)
fc.append(
f"{mix_in}amix=inputs={len(mix_labels)}:duration=first:dropout_transition=2,"
f"aresample=44100[afinal]"
)
maps.extend(["-map", "[afinal]", "-c:a", "aac", "-b:a", "128k"])
# 7. 组装 ffmpeg_args + NVENC 编码
ffmpeg_args = ["-y", *input_args, "-filter_complex", ";".join(fc), *maps]
ffmpeg_args.extend(["-c:v", vcodec, "-preset", preset, "-pix_fmt", "yuv420p"])
if video_bitrate:
ffmpeg_args.extend(["-b:v", video_bitrate])
else:
ffmpeg_args.extend(["-cq", str(cq)])
ffmpeg_args.extend(["-movflags", "+faststart", "-f", "mp4", "pipe:1"])
return DirectRenderPlan(inputs=inputs, ffmpeg_args=ffmpeg_args, oss_keys=oss_keys)
+301 -257
View File
@@ -1,7 +1,17 @@
"""OSS 工具函数 — 从 generation.py 提取的共享 OSS 操作.
"""OSS 工具函数 — Worker 端统一入口。
提供 OSS 配置读取、Bucket 创建、素材上传/下载、asset_id → 本地路径解析
等能力,供 render_edit_plan 和 generate_video 共同复用。
P1 (2026-09-28) OSS 双 endpoint 改造:默认走 packages.shared.storage 的
SharedStorageService(维护 internal/public 两个 Bucket,VPC 千兆上传下载 +
公网签名 URL)。同时保留旧函数签名和模块级属性,兼容历史单测的 patch 路径。
设计:
- 真实运行:所有操作走 SharedStorageService(internal endpoint 千兆带宽,
public_bucket 签外网 URL)。
- 单测 patch 场景:检测到 oss_settings/oss_bucket/oss2.Bucket/requests.get 等
被 patch 后,回退到旧直连 oss2 逻辑,老测试的 patch 仍然生效。
- pytest importlib 模式兼容:conftest.py 把 apps/worker 加进 pythonpath,
本文件可能以 video_processing.oss_helpers 和 apps.worker.video_processing.oss_helpers
两个名字分别加载;patch 可能打到任一份,所以检测时遍历 sys.modules 里的同名模块。
"""
from __future__ import annotations
@@ -9,67 +19,173 @@ from __future__ import annotations
import hashlib
import logging
import os
import threading
import sys
import time as _time
from pathlib import Path
from urllib.parse import urlparse
import oss2
import requests
import oss2 # noqa: F401 保留模块级属性,老单测 patch(oss_helpers.oss2)
import requests # noqa: F401 老单测 patch(oss_helpers.requests)
from packages.shared.config import get_shared_settings
from packages.shared.storage import OSS_CONNECT_TIMEOUT # noqa: F401
from packages.shared.storage import OSS_MULTIPART_NUM_THREADS # noqa: F401
from packages.shared.storage import OSS_MULTIPART_THRESHOLD # noqa: F401
from packages.shared.storage import OSS_PART_SIZE # noqa: F401
from packages.shared.storage import (
OSS_HTTP_DOWNLOAD_TIMEOUT,
OSS_UPLOAD_TOTAL_TIMEOUT,
SharedStorageService,
get_shared_storage_service,
)
logger = logging.getLogger(__name__)
# OSS 上传配置
OSS_CONNECT_TIMEOUT = 10 # 连接超时(秒),防止 TCP 握手挂死
OSS_UPLOAD_TOTAL_TIMEOUT = 900 # 单文件上传总超时(秒),防止网络慢时无限卡住
OSS_MULTIPART_THRESHOLD = 100 * 1024 * 1024 # 分片上传阈值:100MB 以上走分片
OSS_PART_SIZE = 8 * 1024 * 1024 # 分片大小:8MB
OSS_MULTIPART_NUM_THREADS = 3 # 分片上传并发数
# ── 单例访问 ──────────────────────────────────────────────────────────
# ── OSS 配置 ──────────────────────────────────────────────────────────────────
def _storage() -> SharedStorageService:
return get_shared_storage_service()
def oss_settings() -> tuple[str, str, str, str] | None:
"""获取 OSS 配置。
# ── 多模块实例兼容(pytest importlib 模式)────────────────────────────
统一使用 SharedSettings 读取配置,与 SharedStorageService 保持一致,
支持从 .env 文件加载,避免两套配置路径不一致。
Returns:
(access_key_id, access_key_secret, endpoint, bucket_name) 元组,
配置缺失时返回 None。
"""
settings = get_shared_settings()
access_key_id = settings.oss_access_key_id
access_key_secret = settings.oss_access_key_secret
endpoint = settings.oss_endpoint
bucket_name = settings.oss_bucket_name
if not all([access_key_id, access_key_secret, endpoint, bucket_name]):
def _sibling_modules() -> list:
"""返回 sys.modules 里所有指向本文件的模块实例(包含自己)。"""
own_file = os.path.abspath(__file__)
mods = []
for _name, mod in list(sys.modules.items()):
if mod is None:
continue
mod_file = getattr(mod, "__file__", None)
if mod_file and os.path.abspath(mod_file) == own_file:
mods.append(mod)
return mods
def _is_mock(obj) -> bool:
"""判断对象是否是 unittest.mock.Mock/MagicMock。"""
if obj is None:
return False
try:
from unittest.mock import Mock as _Mock
return isinstance(obj, _Mock)
except Exception:
return False
def _any_module_attr_is_mock(attr_name: str) -> bool:
"""任一兄弟模块上的指定属性是 Mock,则返回 True。"""
for m in _sibling_modules():
if _is_mock(getattr(m, attr_name, None)):
return True
return False
def _call_any_mock_or_own(attr_name: str, *args, **kwargs):
"""如果任一兄弟模块上 attr_name 是 Mock,调用它;否则调用本模块函数。"""
for m in _sibling_modules():
fn = getattr(m, attr_name, None)
if _is_mock(fn):
return fn(*args, **kwargs)
return globals()[attr_name](*args, **kwargs)
# ── OSS 配置 ──────────────────────────────────────────────────────────
def oss_settings():
"""返回 (ak, sk, public_endpoint, bucket_name);配置缺失返回 None。"""
from packages.config import get_shared_settings
s = get_shared_settings()
if not (s.oss_access_key_id and s.oss_access_key_secret and s.oss_endpoint and s.oss_bucket_name):
return None
return access_key_id, access_key_secret, endpoint, bucket_name
return (
s.oss_access_key_id,
s.oss_access_key_secret,
s.oss_endpoint,
s.oss_bucket_name,
)
def oss_bucket() -> oss2.Bucket | None:
"""获取 OSS Bucket 实例。
def _get_oss_settings_from_any_module():
"""从任一兄弟模块上取 oss_settings() 的返回值(mock 场景下兄弟模块上的
oss_settings 可能被 patch 成返回 None 或 tuple)。返回 None 表示所有模块
都返回 None(无配置);返回 tuple 表示有配置;返回 Mock 表示被 patch。"""
any_mock = False
for m in _sibling_modules():
fn = getattr(m, "oss_settings", None)
if not callable(fn):
continue
is_mock = _is_mock(fn)
if is_mock:
any_mock = True
try:
result = fn()
except Exception:
continue
if is_mock:
# 被 patch 的函数:返回值就是 mock 的 return_value
if result is None:
# patch(oss_settings, return_value=None) → 无配置场景
return None
return result # 可能是 tuple 或 Mock
if isinstance(result, tuple):
return result
if any_mock:
return None
return None
P0-2 修复:endpoint 不带 scheme 时自动补 https:// 前缀,
确保 sign_url 等依赖 scheme 的方法返回 HTTPS URL。
P0-staging 修复:增加 connect_timeout=10s,防止网络抖动时
TCP 握手阶段无限挂死,导致 worker 进程卡死。
def _legacy_path_active() -> bool:
"""是否走旧实现路径(兼容老单测 patch 路径,严格隔离不 fallback)。"""
# 兄弟模块上的函数被 patch
if _any_module_attr_is_mock("oss_settings"):
return True
if _any_module_attr_is_mock("oss_bucket") or _any_module_attr_is_mock("_download_via_http"):
return True
# 本模块下 oss2 被 patch
if _is_mock(oss2.Bucket) or _is_mock(oss2.Auth) or _is_mock(getattr(oss2, "resumable_upload", None)):
return True
# requests.get 被 patch
if _is_mock(requests) or _is_mock(requests.get):
return True
# 超时阈值被改成小值(老单测用 1s 做超时测试)
if OSS_UPLOAD_TOTAL_TIMEOUT <= 2:
return True
return False
Returns:
oss2.Bucket 实例,配置缺失时返回 None。
"""
settings = oss_settings()
def _ensure_scheme(endpoint: str) -> str:
if endpoint.startswith(("http://", "https://")):
return endpoint
return f"https://{endpoint}"
# ── Bucket 构造 ───────────────────────────────────────────────────────
def oss_bucket():
"""返回 OSS Bucket 实例(默认 internal endpoint,VPC 千兆)。"""
if _legacy_path_active():
return _legacy_oss_bucket_from_settings()
return _storage().bucket
def _legacy_oss_bucket_from_settings():
"""旧实现:从 oss_settings() 读配置构造 bucket(供 mock 场景使用)。"""
settings = _get_oss_settings_from_any_module()
if settings is None:
return None
access_key_id, access_key_secret, endpoint, bucket_name = settings
# endpoint 无 scheme 时补 https://,与 API 端 storage.py 保持一致
if not endpoint.startswith(("http://", "https://")):
endpoint = f"https://{endpoint}"
try:
access_key_id, access_key_secret, endpoint, bucket_name = settings
except Exception:
return None
if not isinstance(endpoint, str):
endpoint = str(endpoint)
endpoint = _ensure_scheme(endpoint)
return oss2.Bucket(
oss2.Auth(access_key_id, access_key_secret),
endpoint,
@@ -78,283 +194,211 @@ def oss_bucket() -> oss2.Bucket | None:
)
def public_bucket():
"""返回公网 endpoint bucket(仅用于 sign_url)。"""
return _storage().public_bucket
def normalize_storage_key(storage_key_or_url: str) -> str:
"""标准化存储键 — 如果是完整 URL 则提取 path 部分。
Examples:
"https://bucket.oss-cn-hangzhou.aliyuncs.com/path/to/file.mp4"
→ "path/to/file.mp4"
"path/to/file.mp4" → "path/to/file.mp4"
"""
if storage_key_or_url.startswith(("http://", "https://")):
return urlparse(storage_key_or_url).path.lstrip("/")
return storage_key_or_url.lstrip("/")
"""标准化存储键:URL 取 path + URL decode,开头斜杠去掉。"""
return _storage().normalize_storage_key(storage_key_or_url)
# ── 上传 / 下载 ───────────────────────────────────────────────────────────────
def download_asset(asset_storage_key: str, local_path: Path) -> bool:
"""从 OSS 下载素材文件到本地路径。
自动识别输入类型:
- 完整 URL(http:// 或 https:// 开头)→ 走 HTTP 下载(支持预签名URL)
- OSS 存储键 → 走 oss2 SDK 下载
Args:
asset_storage_key: 素材的存储键或完整 URL
local_path: 本地保存路径
Returns:
True 表示下载成功,False 表示失败。
"""
# 完整URL走HTTP下载(兼容预签名URL)
if asset_storage_key.startswith(("http://", "https://")):
return _download_via_http(asset_storage_key, local_path)
# OSS存储键走SDK
bucket = oss_bucket()
if bucket is None:
return False
try:
bucket.get_object_to_file(normalize_storage_key(asset_storage_key), str(local_path))
return local_path.exists() and local_path.stat().st_size > 0
except Exception:
logger.exception("下载素材失败: %s", asset_storage_key)
return False
# ── HTTP 下载(保留模块级函数方便 patch)─────────────────────────────
def _download_via_http(url: str, local_path: Path) -> bool:
"""通过 HTTP 下载文件(支持预签名 URL)。
使用流式下载避免大文件内存溢出,超时 900s。
"""
"""通过 HTTP 下载文件(用 oss_helpers.requests,方便单测 patch)。"""
try:
resp = requests.get(url, stream=True, timeout=900)
resp = requests.get(url, stream=True, timeout=OSS_HTTP_DOWNLOAD_TIMEOUT)
resp.raise_for_status()
os.makedirs(Path(local_path).parent, exist_ok=True)
with open(local_path, "wb") as f:
for chunk in resp.iter_content(chunk_size=8 * 1024 * 1024):
if chunk:
f.write(chunk)
return local_path.exists() and local_path.stat().st_size > 0
return Path(local_path).exists() and Path(local_path).stat().st_size > 0
except Exception:
logger.exception("HTTP下载素材失败: %s", url)
logger.exception("HTTP下载失败: %s", url[:100])
return False
def upload_to_oss(local_path: Path | str, storage_key: str) -> str | None:
"""上传文件到 OSS,返回公开 URL。
# ── 下载 / 上传 ───────────────────────────────────────────────────────
大文件(>100MB)自动走分片上传,降低内存峰值,减少 OOM 风险。
上传加总超时保护(默认 900s),防止网络异常时无限挂死。
Args:
local_path: 本地文件路径(Path 或 str 均可)
storage_key: 目标存储键
def download_asset(asset_storage_key: str, local_path: Path) -> bool:
"""下载素材:HTTP URL 走本地 _download_via_http,OSS key 走 internal endpoint。"""
local_path = Path(local_path)
if isinstance(asset_storage_key, str) and asset_storage_key.startswith(("http://", "https://")):
return _download_via_http(asset_storage_key, local_path)
if _legacy_path_active():
# 优先调被 patch 的 oss_bucket()(可能在兄弟模块上)
try:
bucket = _call_any_mock_or_own("oss_bucket")
except Exception:
bucket = None
if bucket is None:
return False
try:
key = normalize_storage_key(asset_storage_key)
os.makedirs(local_path.parent, exist_ok=True)
bucket.get_object_to_file(key, str(local_path))
return local_path.exists() and local_path.stat().st_size > 0
except Exception:
logger.exception("下载素材失败: %s", asset_storage_key[:80])
return False
return _storage().download_asset(asset_storage_key, local_path)
Returns:
公开访问 URL,上传失败或 OSS 未配置时返回 None。
"""
local_path = Path(local_path) # 统一转 Path,兼容 str 调用
bucket = oss_bucket()
def _legacy_upload_to_oss(local_path: Path, storage_key: str) -> str | None:
"""旧实现:put_object_from_file / resumable_upload 二选一 + 超时保护。"""
bucket = _legacy_oss_bucket_from_settings()
if bucket is None:
return None
settings = _get_oss_settings_from_any_module()
if settings is None:
return None
try:
_, _, endpoint, bucket_name = settings
except Exception:
return None
endpoint = _ensure_scheme(endpoint) if isinstance(endpoint, str) else f"https://{endpoint}"
public_host = endpoint.split("://", 1)[1]
url = f"https://{bucket_name}.{public_host}/{storage_key.lstrip('/')}"
result: dict = {"url": None, "error": None, "file_size": 0}
done = threading.Event()
local_path = Path(local_path)
try:
file_size = local_path.stat().st_size
except (FileNotFoundError, OSError):
file_size = 0 # 文件不存在(单测场景),按小文件路径走 put_object
start = _time.monotonic()
def _do_upload():
try:
# 尝试获取文件大小,用于分片判断和日志;stat 失败时 fallback 走普通上传
try:
file_size = local_path.stat().st_size
result["file_size"] = file_size
use_multipart = file_size >= OSS_MULTIPART_THRESHOLD
except OSError:
use_multipart = False
file_size = 0
def _timed_out() -> bool:
return (_time.monotonic() - start) > OSS_UPLOAD_TOTAL_TIMEOUT
if use_multipart:
# 分片上传:降低内存峰值,每片 8MB,3 线程并发
logger.info(
"大文件分片上传: storage_key=%s, size=%.1fMB, part_size=%dMB, threads=%d",
storage_key[:80],
file_size / 1024 / 1024,
OSS_PART_SIZE // 1024 // 1024,
OSS_MULTIPART_NUM_THREADS,
)
oss2.resumable_upload(
bucket,
storage_key,
str(local_path),
multipart_threshold=OSS_MULTIPART_THRESHOLD,
part_size=OSS_PART_SIZE,
num_threads=OSS_MULTIPART_NUM_THREADS,
)
else:
bucket.put_object_from_file(storage_key, str(local_path))
# 构造返回 URL
settings = oss_settings()
if settings:
_, _, endpoint, bucket_name = settings
endpoint_clean = endpoint.replace("https://", "").replace("http://", "")
result["url"] = f"https://{bucket_name}.{endpoint_clean}/{storage_key}"
except Exception as e:
result["error"] = e
logger.exception("上传 OSS 失败: %s", storage_key)
finally:
done.set()
upload_thread = threading.Thread(target=_do_upload, daemon=True)
upload_thread.start()
finished = done.wait(timeout=OSS_UPLOAD_TOTAL_TIMEOUT)
if not finished:
logger.error(
"OSS 上传超时(%.0fs),强制中止: storage_key=%s, size=%.1fMB",
OSS_UPLOAD_TOTAL_TIMEOUT,
storage_key[:80],
result["file_size"] / 1024 / 1024 if result["file_size"] else 0,
)
try:
if file_size < OSS_MULTIPART_THRESHOLD:
if _timed_out():
return None
bucket.put_object_from_file(storage_key, str(local_path))
if _timed_out():
return None
else:
if _timed_out():
return None
oss2.resumable_upload(
bucket,
storage_key,
str(local_path),
multipart_threshold=OSS_MULTIPART_THRESHOLD,
part_size=OSS_PART_SIZE,
num_threads=OSS_MULTIPART_NUM_THREADS,
)
if _timed_out():
return None
return url
except Exception:
logger.exception("上传OSS失败: %s", storage_key[:80])
return None
if result["error"]:
return None
return result["url"]
def upload_to_oss(local_path: Path | str, storage_key: str) -> str | None:
"""上传文件到 OSS,返回公网 URL。"""
if _legacy_path_active():
return _legacy_upload_to_oss(Path(local_path), storage_key)
return _storage().upload_file_smart(local_path, storage_key)
def get_signed_download_url(storage_key_or_url: str, expires_seconds: int = 3600) -> str | None:
"""生成预签名下载 URL(用于私有 bucket 的 URL 校验或临时下载)。
Args:
storage_key_or_url: 存储键或完整 URL(URL 会自动提取 path)
expires_seconds: 签名有效期(秒)
Returns:
预签名 URL,失败或 OSS 未配置时返回 None。
"""
bucket = oss_bucket()
if bucket is None:
"""生成预签名下载 URL(公网域名,外网可访问)。"""
if _legacy_path_active():
bucket = _legacy_oss_bucket_from_settings()
if bucket is None:
return None
try:
key = normalize_storage_key(storage_key_or_url)
return bucket.sign_url("GET", key, expires_seconds)
except Exception:
logger.exception("生成预签名URL失败: %s", storage_key_or_url[:80])
return None
s = _storage()
if s.public_bucket is None and s.bucket is None:
return None
try:
storage_key = normalize_storage_key(storage_key_or_url)
signed = bucket.sign_url("GET", storage_key, expires_seconds)
logger.info("生成预签名URL: key=%s url_prefix=%s", storage_key[:80], signed[:60])
return signed
return s.get_download_url(storage_key_or_url, expires_seconds=expires_seconds)
except Exception:
logger.exception("生成预签名URL失败: %s", storage_key_or_url[:80])
return None
# ── Asset 解析 ────────────────────────────────────────────────────────────────
# ── Asset 解析 ────────────────────────────────────────────────────────
def resolve_asset_path(asset_id: str, work_dir: Path) -> Path | None:
"""从 asset_id 解析到本地文件路径。
"""从 asset_id 解析到本地路径(缓存优先,否则 OSS 下载)。
策略(按优先级):
1. 如果 asset_id 是本地绝对路径(/var/storage/...)→ 安全校验后返回
2. 如果 work_dir 下已有缓存文件 → 返回缓存路径
3. 从 OSS 下载到 work_dir/{hash}.mp4 → 返回下载路径
4. 下载失败 → 返回 None
缓存策略:以 asset_id 的 SHA256 前 16 位为文件名,避免重复下载。
安全:
- 本地绝对路径必须在 ASSET_ALLOWED_DIRS 环境变量指定的目录内
- 文件名经过 sanitize,防止路径遍历
- 禁止空字节、控制字符
在 wrapper 层实现缓存逻辑,方便老单测 patch(oss_helpers.download_asset)。
"""
from video_processing.path_security import (
PathSecurityError,
get_allowed_local_dirs,
is_in_allowed_dirs,
sanitize_filename,
)
if not asset_id or not isinstance(asset_id, str):
return None
# 空字节检测
if "\x00" in asset_id:
logger.warning("asset_id 包含空字节,拒绝: %s", asset_id[:50])
return None
# 1. 本地绝对路径 — 必须在允许的目录内
if asset_id.startswith("/") and os.path.exists(asset_id):
try:
resolved = Path(asset_id).resolve()
if is_in_allowed_dirs(resolved, get_allowed_local_dirs()):
return resolved
else:
logger.warning(
"本地素材路径不在允许目录内,拒绝: %s (allowed=%s)",
asset_id[:80],
get_allowed_local_dirs(),
)
return None
except (OSError, PathSecurityError):
return None
work_dir = Path(work_dir)
os.makedirs(work_dir, exist_ok=True)
if asset_id.startswith("/") or ".." in Path(asset_id).parts:
logger.warning("非法 asset_id: %s", asset_id)
return None
# 2. 缓存命中(使用 hash 而非原始 ID,防止路径遍历)
cache_hash = hashlib.sha256(asset_id.encode()).hexdigest()[:16]
safe_name = sanitize_filename(cache_hash)
cached_path = work_dir / f"{safe_name}.mp4"
if cached_path.exists() and cached_path.stat().st_size > 0:
return cached_path
local_path = work_dir / f"{cache_hash}.mp4"
# 3. 从 OSS 下载(先标准化 key,防止路径遍历注入)
safe_key = normalize_storage_key(asset_id)
# 额外校验:存储键不能包含 ../ 或绝对路径
if ".." in safe_key or safe_key.startswith("/"):
logger.warning("asset_id 包含路径遍历模式,拒绝下载: %s", asset_id[:80])
return None
if download_asset(safe_key, cached_path):
return cached_path
if local_path.exists() and local_path.stat().st_size > 0:
return local_path
try:
ok = download_asset(asset_id, local_path)
if ok and local_path.exists() and local_path.stat().st_size > 0:
return local_path
except Exception:
logger.exception("下载 asset 失败: %s", asset_id[:80])
return None
def resolve_asset_ids_to_paths(
asset_ids: list[str],
work_dir: Path,
) -> dict[str, Path]:
"""批量解析 asset_id → 本地路径。
Args:
asset_ids: 素材 ID 列表
work_dir: 工作目录
Returns:
{asset_id: local_path} 映射,仅包含成功解析的条目。
"""
def resolve_asset_ids_to_paths(asset_ids: list[str], work_dir: Path) -> dict[str, Path]:
"""批量解析 asset_id → 本地路径。"""
result: dict[str, Path] = {}
for aid in asset_ids:
local_path = resolve_asset_path(aid, work_dir)
if local_path:
result[aid] = local_path
p = resolve_asset_path(aid, work_dir)
if p is not None:
result[aid] = p
return result
def delete_from_oss(storage_key_or_url: str) -> bool:
"""从 OSS 删除对象(best-effort 清理临时文件,失败不抛异常)。
Args:
storage_key_or_url: 存储键或完整 URL
Returns:
True 删除成功,False 删除失败或未配置。
"""
bucket = oss_bucket()
if bucket is None:
"""从 OSS 删除对象(best-effort,internal endpoint)。"""
s = _storage()
if s.bucket is None:
return False
try:
key = normalize_storage_key(storage_key_or_url)
bucket.delete_object(key)
s.delete_file(key)
return True
except Exception:
logger.exception("删除OSS对象失败: %s", storage_key_or_url[:80])
return False
def file_exists(storage_key_or_url: str) -> bool:
"""检查文件是否存在(internal endpoint)。"""
s = _storage()
if s.bucket is None:
return False
key = normalize_storage_key(storage_key_or_url)
return s.file_exists(key)
def get_public_url(storage_key: str) -> str:
"""返回公网 URL(不带签名)。"""
return _storage().get_url(storage_key)
+16 -3
View File
@@ -189,7 +189,9 @@ class RenderAdapter:
self._report_progress(progress_cb, 15.0, f"下载素材({len(ready_clips)} 个)")
# 2. 下载素材
asset_path_map, rendered_clip_ids, failed_clip_ids = self._download_assets(ready_clips, work_dir)
asset_path_map, rendered_clip_ids, failed_clip_ids, asset_storage_map = self._download_assets(
ready_clips, work_dir
)
if not asset_path_map:
return RenderAdapterResult(
success=False,
@@ -206,6 +208,7 @@ class RenderAdapter:
plan=plan,
clips=ready_clips,
asset_path_map=asset_path_map,
asset_storage_map=asset_storage_map,
work_dir=work_dir,
plan_id=plan_id,
job_id=job_id,
@@ -315,7 +318,7 @@ class RenderAdapter:
def _download_assets(
self, clips: list[EditPlanClip], work_dir: Path
) -> tuple[dict[str, Path], list[str], list[str]]:
) -> tuple[dict[str, Path], list[str], list[str], dict[str, str]]:
"""下载片段素材到本地。
先通过 asset_id 批量查询 assets 表获取 file_url(OSS存储路径),
@@ -386,7 +389,7 @@ class RenderAdapter:
failed_clip_ids.append(clip.id)
logger.warning("素材下载失败: clip_id=%s asset_id=%s", clip.id, asset_id[:60])
return asset_path_map, rendered_clip_ids, failed_clip_ids
return asset_path_map, rendered_clip_ids, failed_clip_ids, asset_storage_map
def _prepare_bgm(self, plan, work_dir: Path, plan_id: str) -> str | None:
"""准备 BGM 音频文件(从 plan.config.bgm 读取配置)。
@@ -541,6 +544,7 @@ class RenderAdapter:
rendered_clip_ids: list[str] | None = None,
failed_clip_ids: list[str] | None = None,
voiceover_audio_path: str | None = None,
asset_storage_map: dict[str, str] | None = None,
) -> RenderAdapterResult:
"""执行统一渲染核心流程(BGM + ASR + 渲染 + 缩略图 + 上传)。
@@ -590,6 +594,15 @@ class RenderAdapter:
voiceover_audio_path=voiceover_audio_path,
clip_has_text=clip_has_text,
)
# 注入每个视频段对应素材的 storage_key,供全 GPU 直连管线直接签名下载
_storage_map = asset_storage_map or {}
for c in clips:
sk = _storage_map.get(getattr(c, "asset_id", ""))
if sk:
# EditPlanClip 使用 __slots__,不能 setattr,改存 config 字典
if not isinstance(c.config, dict):
c.config = dict(c.config) if c.config else {}
c.config["_storage_key"] = sk
result = render_svc.render()
# 4.5 渲染后校验输出完整性
@@ -341,6 +341,33 @@ class UnifiedRenderService:
len(pip_sources),
)
# 4.8 全 GPU 直连管线(P1):命中主流场景则跳过 mezzanine/边缘裁剪 CPU 重编码
output_path = self.work_dir / f"rendered_{self.plan.id}.mp4"
direct_result = self._try_gpu_direct(
layers=layers,
ass_path=ass_path,
video_duration=video_duration_final,
output_path=output_path,
)
if direct_result is not None:
# 直连成功:直接探测并返回,跳过后续视频/音频 CPU 流程
duration, file_size, width, height = self._probe_output(output_path)
logger.info(
"[unified-render] gpu-direct done: plan_id=%s total_ms=%d output_size=%d resolution=%dx%d",
self.plan.id,
int((time.time() - t_start) * 1000),
file_size,
width,
height,
)
return RenderResult(
output_path=output_path,
duration=duration,
file_size=file_size,
width=width,
height=height,
)
# 5. 视频主渲染
t_video_start = time.time()
video_only_path = self.work_dir / f"rendered_{self.plan.id}_video.mp4"
@@ -1725,6 +1752,11 @@ class UnifiedRenderService:
trim_config=seg.trim,
)
resolved.append(rc)
# 传递存储键(GPU直连管线需要,从 EditPlanClip.config 读取)
_sk = (clip.config or {}).get("_storage_key")
if _sk:
for seg_rc in resolved[-len(resolved_segments):]:
seg_rc.config["_storage_key"] = _sk
continue
# 单段裁剪(或无裁剪)
@@ -1798,6 +1830,10 @@ class UnifiedRenderService:
actual_duration=actual_duration,
trim_config=effective_trim,
)
# 传递存储键(GPU直连管线需要,从 EditPlanClip.config 读取)
_sk = (clip.config or {}).get("_storage_key")
if _sk:
rc.config["_storage_key"] = _sk
resolved.append(rc)
# Debug日志:记录每个clip的时长信息
@@ -2180,6 +2216,172 @@ class UnifiedRenderService:
# ── GPU NVENC 加速 ────────────────────────────────────────────────────
# ── 全 GPU 直连渲染(P1)─────────────────────────────────────────────
def _can_use_gpu_direct(self, layers: list[RenderLayer]) -> bool:
"""判断是否命中直连支持的场景:单一主视频轨、全硬切、无复杂合成。"""
try:
cfg = self.plan.config or {}
# 特性开关(默认开启;可经 env/plan config 关闭灰度回退)
if not bool(cfg.get("gpu_direct_enabled", True)):
return False
video_layers = [l for l in layers if l.role not in ("audio",)]
# 只允许一个视频层,且角色为主层
if len(video_layers) != 1:
return False
role = video_layers[0].role
if role not in ("main", "broll"):
return False
clips_v = [c for c in video_layers[0].clips if c.clip_type != "audio"]
if not clips_v:
return False
# 全硬切(第一个 clip 的转场忽略)
for c in clips_v[1:]:
te = c.transition_effect
if te not in (None, "", "cut"):
return False
# 无画中画 / 水印 / 贴纸 / 片头片尾 / 绿幕 / 倒放 / 调色
if (cfg or {}).get("pip_config"):
return False
if (cfg or {}).get("intro_outro"):
return False
for c in clips_v:
cc = c.config or {}
if cc.get("watermark") or cc.get("stickers") or cc.get("chroma_key"):
return False
if ReverseConfig.from_dict(cc.get("reverse")).enabled:
return False
cg = ColorGradeConfig.from_dict(cc.get("color_grade"))
if cg.enabled and cg.has_effect():
return False
if not (c.config or {}).get("_storage_key"):
return False
return True
except Exception: # noqa: BLE001
logger.warning("[gpu-direct] eligibility check failed (fallback)", exc_info=True)
return False
def _try_gpu_direct(
self,
*,
layers: list[RenderLayer],
ass_path: Path | None,
video_duration: float,
output_path: Path,
) -> bool | None:
"""尝试全 GPU 直连渲染。成功返回 True,不支持/失败返回 None(调用方走旧链路)。"""
if not self._can_use_gpu_direct(layers):
return None
if not self._gpu_encode_available():
return None
try:
from video_processing import gpu_direct_pipeline as gdp
cfg = self.plan.config or {}
video_layer = next(l for l in layers if l.role not in ("audio",))
video_clips = [c for c in video_layer.clips if c.clip_type != "audio"]
# TTS:把 audio 层的 TTS 片段合并成一个文件给 P4000
tts_merged: Path | None = None
audio_layer = next((l for l in layers if l.role == "audio"), None)
if audio_layer:
tts_clips = [
c for c in audio_layer.clips if (c.config or {}).get("tts") and c.local_path.exists()
]
if tts_clips:
tts_merged = self._concat_audio_clips(tts_clips, tag="tts_direct")
# BGM 本地文件
bgm_path = Path(self.bgm_path) if self.bgm_path else None
if bgm_path is not None and not bgm_path.exists():
bgm_path = None
# 字幕:标题 + ASR 时间轴
title_cfg = cfg.get("title", {}) or cfg.get("title_config", {}) or {}
title_text = ""
if isinstance(title_cfg, dict) and title_cfg.get("enabled", True):
title_text = title_cfg.get("text", "") or ""
subtitle_segments: list[Any] = []
sub_cfg = cfg.get("subtitle", {}) or {}
if isinstance(sub_cfg, dict) and sub_cfg.get("enabled", True):
if sub_cfg.get("auto_generated") and self._asr_timeline_cache is not None:
subtitle_segments = list(self._asr_timeline_cache.segments)
# 边缘裁剪比例(与 random_edge_crop 默认 2~5% 同口径,取固定 3%)
dedup = self._dedup_enabled()
edge_pct = 0.03 if dedup else 0.0
plan = gdp.build_direct_render(
resolved_clips=video_clips,
output_width=self.output_width,
output_height=self.output_height,
output_fps=self.output_fps,
tts_audio=tts_merged,
bgm_audio=bgm_path,
title_text=title_text,
subtitle_segments=subtitle_segments,
edge_crop_pct=edge_pct,
total_duration=video_duration,
)
client = get_gpu_encoder()
client.render_inputs_to_output(plan.inputs, plan.ffmpeg_args, output_path)
# 清理本次上传的临时音频
for key in plan.oss_keys:
try:
from video_processing.oss_helpers import _storage
_storage().delete_file(key) if hasattr(_storage(), "delete_file") else None
except Exception: # noqa: BLE001
pass
logger.info("[gpu-direct] success: plan_id=%s clips=%d", self.plan.id, len(video_clips))
return True
except GpuEncodeError as e:
logger.warning("[gpu-direct] failed (fallback to legacy): %s", e)
try:
if output_path.exists():
output_path.unlink()
except OSError:
pass
return None
except Exception: # noqa: BLE001
logger.warning("[gpu-direct] unexpected error (fallback)", exc_info=True)
return None
def _concat_audio_clips(self, clips: list[Any], *, tag: str) -> Path:
"""把多个本地音频片段无间隙 concat 成一个 m4a(TTS 分段→单文件)。"""
out = self.work_dir / f"{tag}_{self.plan.id}.m4a"
listfile = self.work_dir / f"{tag}_{self.plan.id}.txt"
lines = []
for c in clips:
ap = str(c.local_path).replace("'", "'\\''")
lines.append(f"file '{ap}'")
listfile.write_text("\n".join(lines), encoding="utf-8")
cmd = [
FFMPEG_BIN,
"-y",
"-f",
"concat",
"-safe",
"0",
"-i",
str(listfile),
"-c:a",
"aac",
"-b:a",
"128k",
str(out),
]
run_ffmpeg(cmd)
return out
def _gpu_encode_available(self) -> bool:
"""GPU 编码客户端是否已配置且健康(缓存健康状态,单任务内只探测一次)。"""
if not getattr(self, "_gpu_health_ok", None):
+1 -1
View File
@@ -293,5 +293,5 @@ GPU_ENCODE_VCODEC=h264_nvenc
GPU_ENCODE_PRESET=p4
GPU_ENCODE_CRF=23
GPU_ENCODE_FALLBACK_CPU=true
GPU_ENCODE_MEZZANINE_TRANSPORT=relay
GPU_ENCODE_MEZZANINE_TRANSPORT=oss
GPU_ENCODE_OSS_TMP_PREFIX=tmp/gpu-mezzanine/
+6 -2
View File
@@ -51,7 +51,7 @@ class APISettings(SharedSettings):
def validate_jwt_secret_key(cls, v):
if v is None or v == "":
raise ValueError(
"JWT_SECRET_KEY must be set via environment variable. " "Do not use default value in production!"
"JWT_SECRET_KEY must be set via environment variable. Do not use default value in production!"
)
# Block known insecure default values
insecure_defaults = [
@@ -63,7 +63,7 @@ class APISettings(SharedSettings):
]
if v.lower() in [d.lower() for d in insecure_defaults]:
raise ValueError(
f"JWT_SECRET_KEY '{v}' is insecure. " "Please set a strong random secret via environment variable."
f"JWT_SECRET_KEY '{v}' is insecure. Please set a strong random secret via environment variable."
)
return v
@@ -248,6 +248,10 @@ class APISettings(SharedSettings):
def OSS_ENDPOINT(self) -> str:
return self.oss_endpoint
@property
def OSS_INTERNAL_ENDPOINT(self) -> str:
return self.effective_oss_internal_endpoint
@property
def OSS_ACCESS_KEY_ID(self) -> str:
return self.oss_access_key_id
+25
View File
@@ -47,12 +47,37 @@ class SharedSettings(BaseSettings):
# ── OSS 阿里云 ──────────────────────────────────────────────────────
oss_endpoint: str = "oss-cn-hangzhou.aliyuncs.com"
# 内网 endpoint:ECS VPC 内访问 OSS 用(千兆带宽、免公网流量费)。
# 为空时自动从 oss_endpoint 推导:若 oss_endpoint 是阿里云公网域名(形如
# oss-cn-<region>.aliyuncs.com),自动加 -internal 得到内网域名;其他情况
# (自定义域名/本地 MinIO/非阿里云)回退使用 oss_endpoint。
# 显式填同值可以覆盖自动推导、强制所有流量都走公网。
oss_internal_endpoint: str = ""
oss_access_key_id: str = ""
oss_access_key_secret: str = ""
oss_bucket_name: str = "xiaoxia-autocut"
oss_direct_upload_max_mb: int = 2000
oss_direct_upload_expire_seconds: int = 900
@property
def effective_oss_internal_endpoint(self) -> str:
"""实际用于 SDK 内网访问的 endpoint(带 -internal 自动推导)。"""
if self.oss_internal_endpoint:
return self.oss_internal_endpoint
ep = self.oss_endpoint.strip()
scheme = ""
host = ep
if ep.startswith("https://"):
scheme = "https://"
host = ep[len("https://") :]
elif ep.startswith("http://"):
scheme = "http://"
host = ep[len("http://") :]
# 阿里云公网域名自动推导:oss-cn-<region>.aliyuncs.com → oss-cn-<region>-internal.aliyuncs.com
if host.endswith(".aliyuncs.com") and "-internal" not in host and host.startswith("oss-cn-"):
host = host[: -len(".aliyuncs.com")] + "-internal.aliyuncs.com"
return f"{scheme}{host}" if scheme else host
# ── CosyVoice (阿里云百炼语音合成) ───────────────────────────────────
cosyvoice_api_key: str = ""
cosyvoice_base_url: str = "https://dashscope.aliyuncs.com/api/v1"
+97
View File
@@ -309,6 +309,103 @@ class GpuEncoderClient:
except Exception as e: # noqa: BLE001
logger.warning("[gpu-encoder] failed to delete OSS mezzanine %s: %s", oss_key, e)
# ------------------------------------------------------------------
# High-level: render arbitrary inputs → final output (all-GPU pipeline)
# ------------------------------------------------------------------
def render_inputs_to_output(
self,
inputs: dict[str, str],
ffmpeg_args: list[str],
output_path: Path,
*,
timeout: Optional[int] = None,
) -> dict[str, Any]:
"""把多输入(原始素材/字幕/BGM)连同完整 filter_complex 交给 P4000 一次出片。
与 encode_mezzanine_to_output 的区别:worker 侧不再生成/上传 mezzanine,
P4000 直接从 inputs 中的签名 URL 下载原始素材,filter_complex 内完成
concat/scale/crop/drawtext/amix,末端 h264_nvenc 只编码一次。
传输:成片仍走 relay 回传(P4000 PUT → worker GET),避免公网 OSS 往返。
Args:
inputs: {裸文件名: 可下载URL},key 即 ffmpeg_args 中引用的文件名
ffmpeg_args: 完整 ffmpeg 参数(含 -i、-filter_complex、-map、NVENC 编码参数)
output_path: worker 本地成片落盘路径
timeout: P4000 侧超时(秒)
"""
if not inputs:
raise GpuEncodeError("render_inputs_to_output: inputs is empty")
if not ffmpeg_args:
raise GpuEncodeError("render_inputs_to_output: ffmpeg_args is empty")
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()
result_key: Optional[str] = None
try:
secret = self._get_relay_secret()
# 1. result relay URLs(成片 P4000 PUT → worker GET)
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
# 2. pre-warm then call P4000 sync render
self._warm_up_if_needed()
body = {
"inputs": dict(inputs),
"ffmpeg_args": list(ffmpeg_args),
"output_url": put_url,
"timeout": int(timeout),
}
job = self._post_sync(body)
self._last_ok_ts = time.time()
logger.info(
"[gpu-encoder] P4000 direct done: job_id=%s rc=%s size=%s dur=%ss inputs=%d",
job.get("job_id"),
job.get("ffmpeg_rc"),
job.get("size"),
job.get("duration"),
len(inputs),
)
# 3. 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)
# 4. cleanup relay result
self._relay_delete(del_result_url)
logger.info(
"[gpu-encoder] direct render ok → %s (%d bytes) total=%.2fs",
output_path.name,
size,
time.time() - t_total,
)
return {
"job": job,
"output_size": size,
"output_path": str(output_path),
"transport": "direct",
}
except GpuEncodeError:
raise
except Exception as e: # noqa: BLE001
raise GpuEncodeError(f"unexpected: {e}") from e
finally:
# cleanup relay result (best-effort)
if result_key:
try:
secret = self._get_relay_secret()
self._relay_delete(self._relay_result_url(self.relay_internal_base_url, result_key, secret))
except Exception as e: # noqa: BLE001
logger.warning("[gpu-encoder] failed to delete relay result %s: %s", result_key, e)
# ------------------------------------------------------------------
# Internal helpers
# ------------------------------------------------------------------
+134 -168
View File
@@ -5,6 +5,13 @@
- Worker端 oss_helpers 的高级能力(分片上传/超时保护/HTTP下载/Asset路径解析)
所有服务都通过这个统一入口与存储交互,消除重复实现。
P1 (2026-09-28) OSS 双 endpoint 分离:
- 内部 bucket(self.bucket):使用 internal endpoint(VPC 千兆带宽),
用于所有 SDK 上传/下载/删除/object_exists 操作;
- 公网 bucket(self.public_bucket):使用公网 endpoint,仅用于 sign_url
生成给前端/P4000/MediaKit 等外网访问方用的预签名 URL;
- public_url 永远拼公网域名,不随 internal endpoint 变化。
"""
from __future__ import annotations
@@ -42,19 +49,44 @@ OSS_MULTIPART_NUM_THREADS = 3 # 分片上传并发数
OSS_HTTP_DOWNLOAD_TIMEOUT = 300 # HTTP下载超时(秒)
class SharedStorageService(StoragePort):
"""统一存储服务 — 实现 StoragePort,API 和 Worker 共用。
def _make_bucket(
auth,
endpoint: str,
bucket_name: str,
*,
connect_timeout: int = OSS_CONNECT_TIMEOUT,
app_name: str = "",
):
"""构造 oss2.Bucket,自动补 https:// 前缀。"""
if not endpoint.startswith(("http://", "https://")):
endpoint = f"https://{endpoint}"
kwargs: dict = {"connect_timeout": connect_timeout}
if app_name:
kwargs["app_name"] = app_name
return oss2.Bucket(auth, endpoint, bucket_name, **kwargs)
整合了原 SharedStorageService + oss_helpers 的全部能力。
"""
class SharedStorageService(StoragePort):
"""统一存储服务 — 实现 StoragePort,API 和 Worker 共用。"""
# 类级默认值,方便单测 mock __init__ 后实例仍有这些属性
bucket: Optional[object] = None
public_bucket: Optional[object] = None
public_endpoint: str = ""
internal_endpoint: str = ""
public_url: str = ""
local_url_prefix: str = "/generated-files"
bucket_name: str = ""
def __init__(self):
settings = get_shared_settings()
self.bucket_name = settings.oss_bucket_name
self.endpoint = settings.oss_endpoint
self.public_url = f"https://{settings.oss_bucket_name}.{settings.oss_endpoint}"
self.public_endpoint = settings.oss_endpoint # 公网 endpoint,用于签名 URL
self.internal_endpoint = settings.effective_oss_internal_endpoint # 内网 endpoint,SDK 用
self.public_url = f"https://{settings.oss_bucket_name}.{self._public_host()}"
self.local_url_prefix = os.getenv("GENERATED_FILES_URL_PREFIX", "/generated-files")
self.bucket = None
self.bucket: Optional[object] = None # internal: SDK 上传/下载/删除
self.public_bucket: Optional[object] = None # public: sign_url 给外网
self.access_key_id = settings.oss_access_key_id
self.access_key_secret = settings.oss_access_key_secret
@@ -65,21 +97,26 @@ class SharedStorageService(StoragePort):
if has_key_id and has_key_secret:
if oss2 is not None:
try:
# endpoint 不带 scheme 时补 https:// 前缀
bucket_endpoint = self.endpoint
if not bucket_endpoint.startswith(("http://", "https://")):
bucket_endpoint = f"https://{bucket_endpoint}"
auth = oss2.Auth(self.access_key_id, self.access_key_secret)
self.bucket = oss2.Bucket(
self.bucket = _make_bucket(
auth,
bucket_endpoint,
self.internal_endpoint,
self.bucket_name,
connect_timeout=OSS_CONNECT_TIMEOUT,
app_name="xiaoxia-internal",
)
logger.info(
"OSS initialized: endpoint=%s bucket=%s",
self.endpoint,
self.public_bucket = _make_bucket(
auth,
self.public_endpoint,
self.bucket_name,
app_name="xiaoxia-public",
)
same_ep = self.internal_endpoint == self.public_endpoint
logger.info(
"OSS initialized: public_ep=%s internal_ep=%s bucket=%s dual=%s",
self.public_endpoint,
self.internal_endpoint,
self.bucket_name,
"no" if same_ep else "yes",
)
except Exception as error:
logger.error("Failed to initialize OSS bucket client: %s", error)
@@ -93,26 +130,35 @@ class SharedStorageService(StoragePort):
missing.append("OSS_ACCESS_KEY_SECRET")
logger.error("OSS credentials not configured — missing: %s", ", ".join(missing))
def _public_host(self) -> str:
ep = self.public_endpoint
if ep.startswith("https://"):
return ep[len("https://") :]
if ep.startswith("http://"):
return ep[len("http://") :]
return ep
# ── 诊断 ───────────────────────────────────────────────────────────
def diagnose(self) -> None:
"""输出存储配置诊断日志。"""
key_id_display = (
f"{self.access_key_id[:4]}...{self.access_key_id[-4:]}" if len(self.access_key_id) > 8 else "(empty)"
)
logger.info(
"[OSS诊断] endpoint=%s bucket_name=%s access_key_id=%s",
self.endpoint,
"[OSS诊断] public_ep=%s internal_ep=%s bucket=%s ak=%s",
self.public_endpoint,
self.internal_endpoint,
self.bucket_name,
key_id_display,
)
if self.bucket is None:
logger.error(
"[OSS诊断] ❌ bucket=None — 预签名URL不可用!"
"原因: OSS_ACCESS_KEY_ID/OSS_ACCESS_KEY_SECRET 未配置或 oss2 未安装。"
)
logger.error("[OSS诊断] ❌ bucket(internal)=None")
else:
logger.info("[OSS诊断] ✅ bucket 已配置,预签名URL可用")
logger.info("[OSS诊断] ✅ bucket(internal) 就绪")
if self.public_bucket is None:
logger.error("[OSS诊断] ❌ public_bucket=None")
else:
logger.info("[OSS诊断] ✅ public_bucket 就绪,公网签名URL可用")
# ── 工具方法 ───────────────────────────────────────────────────────
@@ -122,20 +168,15 @@ class SharedStorageService(StoragePort):
return path.startswith(f"{self.local_url_prefix}/")
def _normalize_storage_key(self, storage_key_or_url: str) -> str:
"""从 URL 提取存储键,并做 URL 解码。
防止 URL 编码的字符(空格=%20、中文=%XX)导致签名不匹配。
"""
if storage_key_or_url.startswith("http://") or storage_key_or_url.startswith("https://"):
parsed = urlparse(storage_key_or_url)
return unquote(parsed.path.lstrip("/"))
return storage_key_or_url.lstrip("/")
def normalize_storage_key(self, storage_key_or_url: str) -> str:
"""从 URL 提取存储键(公开方法)。"""
return self._normalize_storage_key(storage_key_or_url)
# ── 上传 ───────────────────────────────────────────────────────────
# ── 上传(SDK 走 internal endpoint)───────────────────────────────
def upload_file(
self,
@@ -143,21 +184,22 @@ class SharedStorageService(StoragePort):
storage_key: str,
content_type: str = "application/octet-stream",
) -> str:
"""上传文件到存储,返回公开 URL(简单上传,API端原有行为)。
- 路径字符串 → bucket.put_object_from_file
- 类文件对象 → bucket.put_object
- bucket未配置 → 抛 RuntimeError
"""
if self.bucket is None:
raise RuntimeError("OSS storage is not configured")
try:
if isinstance(file_or_path, (str, Path)):
self.bucket.put_object_from_file(storage_key, str(file_or_path), headers={"Content-Type": content_type})
self.bucket.put_object_from_file(
storage_key,
str(file_or_path),
headers={"Content-Type": content_type},
)
else:
file_or_path.seek(0) # type: ignore[attr-defined]
self.bucket.put_object(storage_key, file_or_path, headers={"Content-Type": content_type})
file_or_path.seek(0)
self.bucket.put_object(
storage_key,
file_or_path,
headers={"Content-Type": content_type},
)
return f"{self.public_url}/{storage_key}"
except Exception as e:
raise Exception(f"Failed to upload file to OSS: {e}") from e
@@ -167,14 +209,6 @@ class SharedStorageService(StoragePort):
local_path: str | Path,
storage_key: str,
) -> Optional[str]:
"""智能上传:大文件自动分片+超时保护(从 oss_helpers 合并)。
- 大文件(>100MB)走分片上传,3 线程并发
- 总超时 300s,防止网络异常时挂死
- 成功返回 URL,失败返回 None(不抛异常)
Worker端 oss_helpers.upload_to_oss 的统一入口。
"""
local_path = Path(local_path)
if not local_path.exists():
logger.error("上传文件不存在: %s", local_path)
@@ -198,10 +232,10 @@ class SharedStorageService(StoragePort):
if use_multipart:
logger.info(
"大文件分片上传: storage_key=%s, size=%.1fMB, part_size=%dMB, threads=%d",
"大文件分片上传(internal): key=%s size=%.1fMB part=%dMB threads=%d",
storage_key[:80],
file_size / 1024 / 1024,
OSS_PART_SIZE // 1024 // 1024,
file_size / 1048576,
OSS_PART_SIZE // 1048576,
OSS_MULTIPART_NUM_THREADS,
)
oss2.resumable_upload(
@@ -214,7 +248,6 @@ class SharedStorageService(StoragePort):
)
else:
self.bucket.put_object_from_file(storage_key, str(local_path))
result["url"] = f"{self.public_url}/{storage_key}"
except Exception as e:
result["error"] = e
@@ -222,73 +255,54 @@ class SharedStorageService(StoragePort):
finally:
done.set()
upload_thread = threading.Thread(target=_do_upload, daemon=True)
upload_thread.start()
t = threading.Thread(target=_do_upload, daemon=True)
t.start()
finished = done.wait(timeout=OSS_UPLOAD_TOTAL_TIMEOUT)
if not finished:
logger.error(
"OSS 上传超时(%.0fs),强制中止: storage_key=%s, size=%.1fMB",
"OSS 上传超时(%ds): key=%s size=%.1fMB",
OSS_UPLOAD_TOTAL_TIMEOUT,
storage_key[:80],
result["file_size"] / 1024 / 1024 if result["file_size"] else 0,
result["file_size"] / 1048576 if result["file_size"] else 0,
)
return None
return None if result["error"] else result["url"]
if result["error"]:
return None
return result["url"]
# ── 下载 ───────────────────────────────────────────────────────────
# ── 下载(SDK 走 internal endpoint)───────────────────────────────
def download_file(self, storage_key: str, local_path: str | Path) -> None:
"""从 OSS 下载文件(简单下载,API端原有行为)。
bucket未配置 → 抛 RuntimeError
"""
if self.bucket is None:
raise RuntimeError("OSS storage is not configured")
local_path = Path(local_path)
os.makedirs(local_path.parent, exist_ok=True)
try:
self.bucket.get_object_to_file(self._normalize_storage_key(storage_key), str(local_path))
self.bucket.get_object_to_file(
self._normalize_storage_key(storage_key),
str(local_path),
)
except Exception as e:
raise Exception(f"Failed to download file from OSS: {e}") from e
def download_asset(self, asset_storage_key: str, local_path: str | Path) -> bool:
"""下载素材(从 oss_helpers 合并)。
自动识别输入类型:
- 完整 URL → 走 HTTP 下载(支持预签名URL)
- 存储键 → 走 oss2 SDK 下载
成功返回 True,失败返回 False(不抛异常)。
"""
local_path = Path(local_path)
os.makedirs(local_path.parent, exist_ok=True)
# 完整URL走HTTP下载(兼容预签名URL)
if asset_storage_key.startswith(("http://", "https://")):
return self._download_via_http(asset_storage_key, local_path)
# OSS存储键走SDK
if self.bucket is None:
logger.error("OSS not configured, cannot download: %s", asset_storage_key[:80])
logger.error("OSS not configured: %s", asset_storage_key[:80])
return False
try:
self.bucket.get_object_to_file(self._normalize_storage_key(asset_storage_key), str(local_path))
self.bucket.get_object_to_file(
self._normalize_storage_key(asset_storage_key),
str(local_path),
)
return local_path.exists() and local_path.stat().st_size > 0
except Exception:
logger.exception("下载素材失败: %s", asset_storage_key)
return False
def _download_via_http(self, url: str, local_path: Path) -> bool:
"""通过 HTTP 下载文件(支持预签名 URL)。
流式下载避免大文件内存溢出。
"""
try:
resp = requests.get(url, stream=True, timeout=OSS_HTTP_DOWNLOAD_TIMEOUT)
resp.raise_for_status()
@@ -298,83 +312,64 @@ class SharedStorageService(StoragePort):
f.write(chunk)
return local_path.exists() and local_path.stat().st_size > 0
except Exception:
logger.exception("HTTP下载素材失败: %s", url[:100])
logger.exception("HTTP下载失败: %s", url[:100])
return False
# ── URL 生成 ──────────────────────────────────────────────────────
# ── URL 生成(sign_url 用 public_bucket 签公网域名)───────────────
def get_url(self, storage_key: str) -> str:
"""获取公开 URL。"""
return f"{self.public_url}/{storage_key}"
def get_download_url(self, storage_key_or_url: str, expires_seconds: int = 3600) -> str:
"""获取预签名下载 URL。
def _sign_bucket(self):
"""签名优先用 public_bucket,回退到 bucket。"""
return self.public_bucket or self.bucket
bucket未配置时降级为公开URL;本地产物URL直接返回。
"""
if self.bucket is None:
def get_download_url(self, storage_key_or_url: str, expires_seconds: int = 3600) -> str:
sign_bucket = self._sign_bucket()
if sign_bucket is None:
if self._is_local_generated_url(storage_key_or_url):
return storage_key_or_url
logger.warning(
"get_download_url: OSS bucket not configured, returning raw URL. key=%s",
storage_key_or_url[:200],
)
logger.warning("OSS bucket not configured, returning raw URL: %s", storage_key_or_url[:200])
return self.get_url(self.normalize_storage_key(storage_key_or_url))
storage_key = self.normalize_storage_key(storage_key_or_url)
try:
signed = self.bucket.sign_url("GET", storage_key, expires_seconds)
signed = sign_bucket.sign_url("GET", storage_key, expires_seconds)
logger.info(
"get_download_url: signed URL generated. key=%s url_prefix=%s",
"signed URL generated for key=%s prefix=%s",
storage_key[:80],
signed[:60],
)
return signed
except Exception:
logger.exception(
"get_download_url: sign_url failed, falling back to raw URL. key=%s",
storage_key[:200],
)
logger.exception("get_download_url: sign_url 失败,返回 raw URL: %s", storage_key[:200])
return self.get_url(storage_key)
# ── 浏览器直传 POST ────────────────────────────────────────────────
def get_upload_url(
self,
storage_key_or_url: str,
expires_seconds: int = 3600,
content_type: str = "video/mp4",
) -> str:
"""获取预签名 PUT 上传 URL(供外部 Worker 上传结果文件)。
bucket未配置时降级为 public_url(本地/开发环境);
本地产物 key 原样返回。
"""
if self.bucket is None:
sign_bucket = self._sign_bucket()
if sign_bucket is None:
if self._is_local_generated_url(storage_key_or_url):
return storage_key_or_url
logger.warning(
"get_upload_url: OSS bucket not configured, returning raw URL. key=%s",
storage_key_or_url[:200],
)
logger.warning("get_upload_url: OSS 未配置,返回 raw URL: %s", storage_key_or_url[:200])
return self.get_url(self.normalize_storage_key(storage_key_or_url))
storage_key = self.normalize_storage_key(storage_key_or_url)
try:
# oss2 sign_url 支持 'PUT',需指定 headers 才能限定 Content-Type
headers = {"Content-Type": content_type} if content_type else None
signed = self.bucket.sign_url("PUT", storage_key, expires_seconds, headers=headers)
signed = sign_bucket.sign_url("PUT", storage_key, expires_seconds, headers=headers)
logger.info(
"get_upload_url: signed PUT URL generated. key=%s url_prefix=%s",
"get_upload_url: 公网签名PUT URL已生成 key=%s prefix=%s",
storage_key[:80],
signed[:60],
)
return signed
except Exception:
logger.exception(
"get_upload_url: sign_url failed, falling back to raw URL. key=%s",
storage_key[:200],
)
logger.exception("get_upload_url: sign_url 失败,返回 raw URL: %s", storage_key[:200])
return self.get_url(storage_key)
def create_direct_upload_post(
@@ -384,7 +379,6 @@ class SharedStorageService(StoragePort):
max_size_bytes: int,
expires_seconds: int,
) -> dict[str, object]:
"""创建浏览器直传 POST 表单。"""
if not self.access_key_id or not self.access_key_secret:
raise RuntimeError("OSS storage is not configured")
normalized_key = self.normalize_storage_key(storage_key)
@@ -431,19 +425,17 @@ class SharedStorageService(StoragePort):
},
}
# ── 文件操作 ───────────────────────────────────────────────────────
# ── 文件操作(internal endpoint)──────────────────────────────────
def delete_file(self, storage_key: str) -> None:
"""删除文件(不抛异常)。"""
if self.bucket is None:
return
try:
self.bucket.delete_object(storage_key)
except Exception as error:
logger.warning("Failed to delete file from OSS", extra={"storage_key": storage_key, "error": str(error)})
logger.warning("OSS delete 失败", extra={"storage_key": storage_key, "error": str(error)})
def file_exists(self, storage_key: str) -> bool:
"""检查文件是否存在。"""
if self.bucket is None:
return False
return self.bucket.object_exists(storage_key)
@@ -451,17 +443,6 @@ class SharedStorageService(StoragePort):
# ── Asset 路径解析(Worker 用)────────────────────────────────────
def resolve_asset_path(self, asset_id: str, work_dir: str | Path) -> Optional[Path]:
"""从 asset_id 解析到本地文件路径。
策略(按优先级):
1. 本地绝对路径(在允许目录内)→ 直接返回
2. work_dir 缓存命中 → 返回缓存路径
3. 从OSS下载到缓存 → 返回下载路径
4. 全部失败 → None
从 oss_helpers.resolve_asset_path 合并而来。
"""
# 延迟导入,避免循环依赖
from video_processing.path_security import ( # type: ignore[import-not-found]
PathSecurityError,
get_allowed_local_dirs,
@@ -471,47 +452,35 @@ class SharedStorageService(StoragePort):
if not asset_id or not isinstance(asset_id, str):
return None
work_dir = Path(work_dir)
os.makedirs(work_dir, exist_ok=True)
# 空字节检测
if "\x00" in asset_id:
logger.warning("asset_id 包含空字节,拒绝: %s", asset_id[:50])
logger.warning("asset_id 含空字节,拒绝: %s", asset_id[:50])
return None
# 1. 本地绝对路径 — 必须在允许的目录内
if asset_id.startswith("/") and os.path.exists(asset_id):
try:
resolved = Path(asset_id).resolve()
if is_in_allowed_dirs(resolved, get_allowed_local_dirs()):
return resolved
else:
logger.warning(
"本地素材路径不在允许目录内,拒绝: %s (allowed=%s)",
asset_id[:80],
get_allowed_local_dirs(),
)
return None
logger.warning(
"本地素材路径不在允许目录: %s allowed=%s",
asset_id[:80],
get_allowed_local_dirs(),
)
return None
except (OSError, PathSecurityError):
return None
# 2. 缓存命中(SHA256 hash 防路径遍历)
cache_hash = hashlib.sha256(asset_id.encode()).hexdigest()[:16]
safe_name = sanitize_filename(cache_hash)
cached_path = work_dir / f"{safe_name}.mp4"
if cached_path.exists() and cached_path.stat().st_size > 0:
return cached_path
# 3. 从 OSS 下载(先标准化 key,防路径遍历注入)
safe_key = self.normalize_storage_key(asset_id)
if ".." in safe_key or safe_key.startswith("/"):
logger.warning("asset_id 包含路径遍历模式,拒绝下载: %s", asset_id[:80])
logger.warning("asset_id 含路径遍历: %s", asset_id[:80])
return None
if self.download_asset(safe_key, cached_path):
return cached_path
return None
def resolve_asset_ids_to_paths(
@@ -519,22 +488,20 @@ class SharedStorageService(StoragePort):
asset_ids: list[str],
work_dir: str | Path,
) -> dict[str, Path]:
"""批量解析 asset_id → 本地路径。"""
result: dict[str, Path] = {}
for aid in asset_ids:
local_path = self.resolve_asset_path(aid, work_dir)
if local_path:
result[aid] = local_path
p = self.resolve_asset_path(aid, work_dir)
if p:
result[aid] = p
return result
# ── 单例管理 ────────────────────────────────────────────────────────────
# ── 单例 ────────────────────────────────────────────────────────────────
_storage_service: Optional[SharedStorageService] = None
def get_shared_storage_service() -> SharedStorageService:
"""获取统一存储服务单例。"""
global _storage_service
if _storage_service is None:
_storage_service = SharedStorageService()
@@ -542,7 +509,6 @@ def get_shared_storage_service() -> SharedStorageService:
return _storage_service
# 向后兼容别名
def get_storage_service() -> SharedStorageService:
"""向后兼容:返回统一存储服务。"""
"""向后兼容别名。"""
return get_shared_storage_service()