Files
xiaoxia-saas/apps/worker/video_processing/oss_helpers.py
T
saas-backend-agent 1a917a3472
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 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 47s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 1m56s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 2m21s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 2m47s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m11s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 3m35s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 6m23s
AI Code Review / AI Code Review (pull_request) Successful in 6m36s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 8m57s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 12m37s
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 - Style (pull_request) Has been cancelled
feat(p1): OSS 双 endpoint 分离 — 上传/下载走 VPC 内网,签名 URL 走公网
问题:
- 116 是阿里云杭州 ECS,之前 oss_endpoint=oss-cn-hangzhou.aliyuncs.com(公网),
  worker 上传/下载 OSS 走公网卡到 0.3~0.7 MB/s,mezzanine 上传 104s + 成片上传 55s,
  吃掉 160s+ 渲染耗时。
- oss_helpers 自己维护一套 bucket 实例,和 SharedStorageService 重复实现,且
  normalize_storage_key 没做 URL decode(unquote),带编码字符的 key 签名出来
  返回 403,触发 _verify_url_accessible 2 次重试延迟。

改造:
1. config/base.py 新增 oss_internal_endpoint 配置项 + effective_oss_internal_endpoint
   属性:阿里云公网域名(oss-cn-<region>.aliyuncs.com)自动推导 -internal 内网
   域名,显式配置可覆盖;非阿里云/MinIO 回退公网。
2. packages/shared/storage.py 维护两个 bucket:
   - self.bucket:internal endpoint(VPC 千兆带宽),SDK 上传/下载/删除/exists
   - self.public_bucket:public endpoint,仅用于 sign_url 生成公网可访问 URL
   public_url 永远拼公网域名,对外不可感知。
3. apps/worker/video_processing/oss_helpers.py 收敛为 SharedStorageService 的
   薄 wrapper(-226 行),保留所有旧函数签名兼容现有调用:
   - oss_bucket() 返回 internal bucket
   - upload_to_oss/download_asset/delete_from_oss 转发到 storage 走 internal
   - get_signed_download_url 转发到 storage.get_download_url 用 public_bucket 签名
   - resolve_asset_path 转发到 storage 统一逻辑
4. normalize_storage_key 统一 unquote,修复签名 403。
5. config/api_settings.py 增加 OSS_INTERNAL_ENDPOINT 属性方便诊断。

预期效果:
- OSS 上传/下载从 0.3~0.7 MB/s → VPC 千兆(~100 MB/s),两段上传 160s 压到 <3s
- 签名 URL 不再 403,_verify_url_accessible 首试即过,不再重试
- 整体渲染耗时目标 <30s(GPU 编码链路 ~10s + OSS 传输 <5s + 合成/封面 <15s)
- 前端/P4000/MediaKit 外网访问 OSS 完全不变(public_url 仍是公网域名)
2026-09-28 09:58:33 +08:00

163 lines
6.2 KiB
Python
Executable File
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""OSS 工具函数 — SharedStorageService 的薄封装。
历史原因 Worker 有一套自己的 oss_helpers 函数(oss_bucket / upload_to_oss /
download_asset / get_signed_download_url / resolve_asset_path 等)。P1
双 endpoint 改造后,所有 OSS 操作统一走 packages.shared.storage 里的
SharedStorageService(维护 internal/public 两个 Bucket),本文件仅保留
Worker 端惯用的函数签名,内部转发到 SharedStorageService。
注意:
- oss_bucket() 返回 internal bucket(用于 SDK 上传/下载/删除/object_exists,
走 VPC 千兆带宽);
- get_signed_download_url() 内部走 SharedStorageService.get_download_url(),
自动用 public_bucket 签公网域名,给 MediaKit/P4000 等外网访问方用;
- upload_to_oss() 走 internal endpoint 上传,返回公网 URL。
"""
from __future__ import annotations
import logging
from pathlib import Path
from packages.shared.storage import (
OSS_CONNECT_TIMEOUT, # noqa: F401 向后兼容
OSS_HTTP_DOWNLOAD_TIMEOUT,
OSS_MULTIPART_NUM_THREADS,
OSS_PART_SIZE, # noqa: F401 向后兼容
OSS_MULTIPART_THRESHOLD,
OSS_UPLOAD_TOTAL_TIMEOUT,
SharedStorageService,
get_shared_storage_service,
)
logger = logging.getLogger(__name__)
# 向后兼容常量(旧代码直接 import OSS_UPLOAD_TOTAL_TIMEOUT 等)
OSS_MULTIPART_NUM_THREADS = OSS_MULTIPART_NUM_THREADS
OSS_MULTIPART_THRESHOLD = OSS_MULTIPART_THRESHOLD
OSS_UPLOAD_TOTAL_TIMEOUT = OSS_UPLOAD_TOTAL_TIMEOUT
OSS_HTTP_DOWNLOAD_TIMEOUT = OSS_HTTP_DOWNLOAD_TIMEOUT
# ── 单例访问 ──────────────────────────────────────────────────────────
def _storage() -> SharedStorageService:
return get_shared_storage_service()
# ── OSS 配置 / Bucket(兼容旧 API)────────────────────────────────────
def oss_settings():
"""兼容旧调用:返回 (ak, sk, public_endpoint, bucket_name)。
注意:endpoint 这里返回公网 endpoint(签名/拼 URL 用),
内部上传/下载 SDK 操作实际走 internal_endpoint(在 SharedStorageService 内)。
"""
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 (
s.oss_access_key_id,
s.oss_access_key_secret,
s.oss_endpoint,
s.oss_bucket_name,
)
def oss_bucket():
"""返回 internal endpoint bucket(用于 SDK 操作:上传/下载/删除/object_exists)。
P1 双 endpoint 改造后,Worker 内所有 SDK 调用应走 VPC internal endpoint,
签名 URL 请改用 get_signed_download_url() 内部自动用 public_bucket。
"""
return _storage().bucket
def public_bucket():
"""返回公网 endpoint bucket(仅用于 sign_url,正常业务不要直接用)。"""
return _storage().public_bucket
def normalize_storage_key(storage_key_or_url: str) -> str:
"""标准化存储键 — URL 提取 path 部分,URL 解码。"""
return _storage().normalize_storage_key(storage_key_or_url)
# ── 上传 / 下载 ──────────────────────────────────────────────────────
def download_asset(asset_storage_key: str, local_path: Path) -> bool:
"""下载素材:URL 走 HTTP,存储键走 OSS SDK(internal endpoint)。"""
return _storage().download_asset(asset_storage_key, local_path)
def upload_to_oss(local_path: Path | str, storage_key: str) -> str | None:
"""上传文件到 OSS(internal endpoint),返回公网 URL。
大文件自动分片(>100MB,3 线程,8MB/片),总超时 300s。
"""
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(公网域名,P4000/MediaKit/浏览器可用)。
内部用 public_bucket 签名,保证签名 URL 的 host 是公网 endpoint。
"""
storage = _storage()
if storage.public_bucket is None and storage.bucket is None:
return None
try:
return storage.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 解析 ────────────────────────────────────────────────────────
def resolve_asset_path(asset_id: str, work_dir: Path) -> Path | None:
"""从 asset_id 解析到本地路径(缓存优先,否则从 OSS 下载)。"""
return _storage().resolve_asset_path(asset_id, work_dir)
def resolve_asset_ids_to_paths(asset_ids: list[str], work_dir: Path) -> dict[str, Path]:
"""批量解析 asset_id → 本地路径。"""
return _storage().resolve_asset_ids_to_paths(asset_ids, work_dir)
def delete_from_oss(storage_key_or_url: str) -> bool:
"""从 OSS 删除对象(internal endpoint,best-effort)。"""
storage = _storage()
if storage.bucket is None:
return False
try:
key = normalize_storage_key(storage_key_or_url)
storage.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)。"""
storage = _storage()
if storage.bucket is None:
return False
key = normalize_storage_key(storage_key_or_url)
return storage.file_exists(key)
def get_public_url(storage_key: str) -> str:
"""返回公网 URL(不带签名,bucket 需公有读或配合签名使用)。"""
return _storage().get_url(storage_key)