Files
xiaoxia-saas/apps/worker/video_processing/oss_helpers.py
T
CI Bot d3fc05d32b
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 Worker Image (pull_request) Successful in 41s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 43s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 1m59s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 2m11s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m39s
CI/CD Pipeline / Validate - Style (pull_request) Successful in 2m49s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m2s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 4m23s
AI Code Review / AI Code Review (pull_request) Successful in 6m37s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 8m38s
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Web 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 5m48s
style: auto-format with black + isort + ruff + prettier [skip ci-format-check]
2026-09-28 02:16:16 +00:00

163 lines
6.3 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 向后兼容
from packages.shared.storage import OSS_PART_SIZE # noqa: F401 向后兼容
from packages.shared.storage import (
OSS_HTTP_DOWNLOAD_TIMEOUT,
OSS_MULTIPART_NUM_THREADS,
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)