Files
xiaoxia-saas/apps/api/app/api/routes/upload.py
T
xiaoxia ec45d71d2c
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (push) Successful in 7s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 8s
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 7s
CI/CD Pipeline / Check push changed paths (push) Successful in 19s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 31s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 52s
CI/CD Pipeline / Build Staging API Image (push) Successful in 45s
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 51s
CI/CD Pipeline / Retag skipped Staging Web Image (push) Successful in 29s
CI/CD Pipeline / Integration Tests (push) Successful in 2m47s
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 2m47s
CI/CD Pipeline / Validate - Style (push) Successful in 3m1s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 3m12s
CI/CD Pipeline / PR Build Web Image (pull_request) Successful in 3m25s
CI/CD Pipeline / CI Gate (pull_request) Successful in 2s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m45s
AI Code Review / AI Code Review (pull_request) Failing after 3m50s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 1m16s
CI/CD Pipeline / ACR Image Cleanup (push) Successful in 1m31s
CI/CD Pipeline / Frontend Unit Tests (push) Successful in 6m28s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 3m24s
CI/CD Pipeline / Validate - Security (push) Successful in 8m20s
CI/CD Pipeline / Staging E2E Tests (push) Successful in 4m16s
CI/CD Pipeline / Unit Tests (push) Successful in 8m54s
CI/CD Pipeline / Build Production Worker Image (push) Successful in 20s
CI/CD Pipeline / Build Production API Image (push) Successful in 22s
CI/CD Pipeline / Build Production Web Image (push) Successful in 2m9s
CI/CD Pipeline / Deploy Production (push) Failing after 11m33s
CI/CD Pipeline / Production Browser E2E (push) Has been skipped
CI/CD Pipeline / Canary Release to Production (push) Failing after 312h17m4s
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Failing after 312h25m56s
CI/CD Pipeline / Retag skipped Staging API Image (push) Failing after 312h25m59s
CI/CD Pipeline / Canary Release to Production (pull_request) Failing after 312h26m5s
CI/CD Pipeline / Staging API Integration Tests (pull_request) Failing after 312h26m21s
CI/CD Pipeline / Staging E2E Tests (pull_request) Failing after 312h26m24s
CI/CD Pipeline / Build Production Worker Image (pull_request) Failing after 312h26m33s
CI/CD Pipeline / Build Production Web Image (pull_request) Failing after 312h26m34s
CI/CD Pipeline / Build Production API Image (pull_request) Failing after 312h26m35s
CI/CD Pipeline / Build Staging Web Image (push) Failing after 312h26m54s
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Failing after 312h27m22s
CI/CD Pipeline / Frontend Unit Tests (pull_request) Failing after 312h28m1s
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Failing after 312h27m22s
CI/CD Pipeline / Frontend Lint (pull_request) Failing after 312h28m1s
CI/CD Pipeline / Validate - Style (pull_request) Failing after 312h28m3s
CI/CD Pipeline / Integration Tests (pull_request) Failing after 312h28m1s
CI/CD Pipeline / Frontend Lint (push) Failing after 312h28m12s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 312h28m2s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Failing after 312h28m14s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Failing after 312h28m2s
CI/CD Pipeline / Build Staging Worker Image (pull_request) Failing after 312h28m14s
CI/CD Pipeline / PR Build Worker Image (push) Failing after 312h28m20s
CI/CD Pipeline / Build Staging Web Image (pull_request) Failing after 312h28m14s
CI/CD Pipeline / PR Build Web Image (push) Failing after 312h28m20s
CI/CD Pipeline / Check push changed paths (pull_request) Failing after 312h28m21s
CI/CD Pipeline / Build Staging API Image (pull_request) Failing after 312h28m15s
CI/CD Pipeline / Check if frontend-only change (push) Failing after 312h28m23s
CI/CD Pipeline / PR Build API Image (push) Failing after 312h28m20s
CI/CD Pipeline / CI Gate (push) Failing after 312h53m37s
CI/CD Pipeline / ACR Image Cleanup (pull_request) Failing after 313h0m42s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Failing after 313h0m53s
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Failing after 313h1m45s
CI/CD Pipeline / Validate - Security (pull_request) Failing after 313h2m25s
fix: 上传素材后视频库立即显示(创建 PROCESSING 状态 Asset) (#1644)
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com>
Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
2026-09-03 14:56:35 +08:00

367 lines
13 KiB
Python

import logging
from typing import Any
from uuid import uuid4
from app.api.routes._helpers import require_project_and_library
from app.auth import AuthenticatedUser, get_current_user
from app.config import get_settings
from app.core.celery_app import celery_app
from app.core.storage import OSSStorageService, get_storage_service
from app.dependencies import (
get_asset_library_repository,
get_asset_repository,
get_ingest_job_repository,
get_project_repository,
)
from app.schemas.upload import (
DirectUploadCompleteRequest,
DirectUploadCompleteResponse,
DirectUploadPrepareRequest,
DirectUploadPrepareResponse,
UploadAssetResponse,
)
from fastapi import APIRouter, Depends, File, Form, HTTPException, UploadFile, status
from packages.application import SubmitIngestJobCommand, SubmitIngestJobUseCase
from packages.domain import Asset, AssetStatus
logger = logging.getLogger(__name__)
router = APIRouter()
# 允许上传的文件 MIME 类型
ALLOWED_MIME_TYPES = frozenset(
{
# 视频
"video/mp4",
"video/mpeg",
"video/quicktime",
"video/x-msvideo",
"video/webm",
"video/x-matroska",
"video/3gpp",
# 音频
"audio/mpeg",
"audio/wav",
"audio/ogg",
"audio/flac",
"audio/aac",
"audio/mp3",
"audio/x-m4a",
"audio/webm",
# 图片
"image/jpeg",
"image/png",
"image/gif",
"image/webp",
"image/bmp",
"image/tiff",
"image/svg+xml",
}
)
def _validate_mime_type(content_type: str | None) -> str:
"""验证并返回标准化的 MIME 类型,如果无效则抛出异常。"""
if not content_type:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail="Content-Type header is required",
)
# 处理带参数的类型,如 "video/mp4; charset=utf-8"
base_type = content_type.split(";")[0].strip().lower()
if base_type not in ALLOWED_MIME_TYPES:
raise HTTPException(
status_code=status.HTTP_415_UNSUPPORTED_MEDIA_TYPE,
detail=f"File type '{base_type}' is not supported. Allowed types: video, audio, and image files.",
)
return base_type
def _infer_mime_type_from_storage_key(storage_key: str) -> str:
"""从 storage_key 推断 MIME 类型(与 worker 端保持一致)。"""
lower_filename = storage_key.rsplit("/", 1)[-1].lower()
_MIME_MAP = {
".mov": "video/quicktime", ".mp4": "video/mp4", ".avi": "video/x-msvideo",
".mkv": "video/x-matroska", ".webm": "video/webm",
".png": "image/png", ".gif": "image/gif", ".bmp": "image/bmp",
".svg": "image/svg+xml", ".jpg": "image/jpeg", ".jpeg": "image/jpeg",
".mp3": "audio/mpeg", ".wav": "audio/wav", ".ogg": "audio/ogg",
".flac": "audio/flac", ".m4a": "audio/x-m4a",
}
for ext, mime in _MIME_MAP.items():
if lower_filename.endswith(ext):
return mime
return "video/mp4" # default
def _create_pending_asset(
asset_repository, project_id, library_id, storage_key, filename, mime_type, user_id, file_hash=""
):
"""立即创建一条 PROCESSING 状态的 Asset 记录,使前端能马上看到新素材。"""
asset = Asset.create(
project_id=project_id,
library_id=library_id,
name=filename,
storage_key=storage_key,
mime_type=mime_type,
status=AssetStatus.PROCESSING,
uploaded_by_user_id=user_id,
file_hash=file_hash,
)
return asset_repository.create(asset)
def _submit_ingest_job(
project_id: str,
library_id: str,
storage_key: str,
ingest_job_repository: Any,
file_hash: str = "",
) -> Any:
use_case = SubmitIngestJobUseCase(ingest_job_repository)
job = use_case.execute(
SubmitIngestJobCommand(
project_id=project_id,
library_id=library_id,
storage_key=storage_key,
file_hash=file_hash,
)
)
celery_app.send_task("worker.ingest_asset", args=[job.id])
return job
@router.post("/direct/prepare", response_model=DirectUploadPrepareResponse)
async def prepare_direct_upload(
request: DirectUploadPrepareRequest,
authenticated_user: AuthenticatedUser = Depends(get_current_user),
project_repository: Any = Depends(get_project_repository),
asset_library_repository: Any = Depends(get_asset_library_repository),
storage_service: OSSStorageService = Depends(get_storage_service),
) -> DirectUploadPrepareResponse:
"""创建浏览器直传 OSS 的短期表单签名。"""
settings = get_settings()
max_size_bytes = settings.OSS_DIRECT_UPLOAD_MAX_MB * 1024 * 1024
if request.file_size > max_size_bytes:
raise HTTPException(
status_code=status.HTTP_413_REQUEST_ENTITY_TOO_LARGE,
detail=f"File exceeds upload limit ({settings.OSS_DIRECT_UPLOAD_MAX_MB}MB)",
)
# P2-5: 服务端验证 MIME 类型
validated_content_type = _validate_mime_type(request.content_type)
require_project_and_library(
request.project_id,
request.library_id,
project_repository,
asset_library_repository,
)
file_id = uuid4().hex[:8]
safe_filename = request.filename.replace("/", "_").replace("\\", "_")
storage_key = f"uploads/{file_id}/{safe_filename}"
try:
payload = storage_service.create_direct_upload_post(
storage_key=storage_key,
content_type=validated_content_type,
max_size_bytes=max_size_bytes,
expires_seconds=settings.OSS_DIRECT_UPLOAD_EXPIRE_SECONDS,
)
except RuntimeError as error:
logger.error("OSS not configured for direct upload prepare: %s", error)
raise HTTPException(status_code=status.HTTP_503_SERVICE_UNAVAILABLE, detail=str(error)) from error
except Exception as error:
logger.exception("Unexpected error in direct upload prepare: %s", error)
raise HTTPException(
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
detail=f"Failed to prepare upload: {type(error).__name__}",
) from error
return DirectUploadPrepareResponse(
upload_url=str(payload["url"]),
method=str(payload["method"]),
storage_key=str(payload["storage_key"]),
expires_at=str(payload["expires_at"]),
fields={str(key): str(value) for key, value in dict(payload["fields"]).items()},
max_size_bytes=max_size_bytes,
)
@router.post("/direct/complete", response_model=DirectUploadCompleteResponse)
async def complete_direct_upload(
request: DirectUploadCompleteRequest,
authenticated_user: AuthenticatedUser = Depends(get_current_user),
ingest_job_repository: Any = Depends(get_ingest_job_repository),
project_repository: Any = Depends(get_project_repository),
asset_library_repository: Any = Depends(get_asset_library_repository),
asset_repository: Any = Depends(get_asset_repository),
storage_service: OSSStorageService = Depends(get_storage_service),
) -> DirectUploadCompleteResponse:
"""确认浏览器直传完成并创建导入任务。"""
require_project_and_library(
request.project_id,
request.library_id,
project_repository,
asset_library_repository,
)
normalized_key = storage_service._normalize_storage_key(request.storage_key)
if not normalized_key.startswith("uploads/"):
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="Invalid upload key")
try:
file_exists = storage_service.file_exists(normalized_key)
except Exception as error:
logger.exception("OSS error checking file existence for key=%s: %s", normalized_key, error)
raise HTTPException(
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
detail="Storage service unavailable",
) from error
if not file_exists:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Uploaded file not found")
# ── 素材去重检测:同素材库 + 同 file_hash 视为重复 ──
if request.file_hash:
existing = asset_repository.find_by_library_and_file_hash(
library_id=request.library_id,
file_hash=request.file_hash,
)
if existing is not None:
logger.info(
"素材去重命中: library=%s hash=%s existing_asset=%s",
request.library_id,
request.file_hash,
existing.id,
)
return DirectUploadCompleteResponse(
storage_key=normalized_key,
ingest_job_id="",
duplicated=True,
asset_id=existing.id,
url=storage_service.get_url(normalized_key),
)
# 立即创建 Asset 记录(PROCESSING 状态),使前端刷新后即可看到新素材
filename = normalized_key.rsplit("/", 1)[-1]
mime_type = _infer_mime_type_from_storage_key(normalized_key)
pending_asset = _create_pending_asset(
asset_repository=asset_repository,
project_id=request.project_id,
library_id=request.library_id,
storage_key=normalized_key,
filename=filename,
mime_type=mime_type,
user_id=authenticated_user.user.id,
file_hash=request.file_hash,
)
job = _submit_ingest_job(
project_id=request.project_id,
library_id=request.library_id,
storage_key=normalized_key,
ingest_job_repository=ingest_job_repository,
file_hash=request.file_hash,
)
return DirectUploadCompleteResponse(
storage_key=normalized_key,
ingest_job_id=job.id,
asset_id=pending_asset.id,
url=storage_service.get_url(normalized_key),
)
@router.post(
"",
response_model=UploadAssetResponse,
summary="Upload Asset",
description="上传素材文件(multipart/form-data),支持视频、音频、图片。触发导入流水线自动处理。",
)
async def upload_asset(
project_id: str = Form(..., min_length=1, description="项目 ID"),
library_id: str = Form(..., min_length=1, description="素材库 ID"),
file: UploadFile = File(..., description="要上传的文件(视频、音频、图片等)"),
file_hash: str = Form(default="", description="文件 MD5 哈希,用于去重检测"),
authenticated_user: AuthenticatedUser = Depends(get_current_user),
ingest_job_repository: Any = Depends(get_ingest_job_repository),
project_repository: Any = Depends(get_project_repository),
asset_library_repository: Any = Depends(get_asset_library_repository),
asset_repository: Any = Depends(get_asset_repository),
storage_service: OSSStorageService = Depends(get_storage_service),
) -> UploadAssetResponse:
"""上传素材文件并触发导入流水线。"""
require_project_and_library(project_id, library_id, project_repository, asset_library_repository)
# ── 素材去重检测:上传前检查同素材库 + 同 file_hash ──
if file_hash:
existing = asset_repository.find_by_library_and_file_hash(
library_id=library_id,
file_hash=file_hash,
)
if existing is not None:
logger.info(
"素材去重命中(multipart): library=%s hash=%s existing_asset=%s",
library_id,
file_hash,
existing.id,
)
return UploadAssetResponse(
storage_key=existing.storage_key,
ingest_job_id="",
url="",
duplicated=True,
asset_id=existing.id,
)
# P2-5: 服务端验证 MIME 类型
validated_content_type = _validate_mime_type(file.content_type)
file_id = uuid4().hex[:8]
safe_filename = file.filename.replace("/", "_").replace("\\", "_") if file.filename else "unknown"
storage_key = f"uploads/{file_id}/{safe_filename}"
try:
file_url = storage_service.upload_file(
file.file,
storage_key,
content_type=validated_content_type,
)
except RuntimeError as error:
logger.error("OSS not configured for upload: %s", error)
raise HTTPException(status_code=status.HTTP_503_SERVICE_UNAVAILABLE, detail=str(error)) from error
except Exception as error:
logger.exception("Unexpected error uploading file to OSS: %s", error)
raise HTTPException(
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
detail=f"Failed to upload file: {type(error).__name__}",
) from error
# 立即创建 Asset 记录(PROCESSING 状态),使前端刷新后即可看到新素材
pending_asset = _create_pending_asset(
asset_repository=asset_repository,
project_id=project_id,
library_id=library_id,
storage_key=storage_key,
filename=safe_filename,
mime_type=validated_content_type,
user_id=authenticated_user.user.id,
file_hash=file_hash,
)
job = _submit_ingest_job(
project_id=project_id,
library_id=library_id,
storage_key=storage_key,
ingest_job_repository=ingest_job_repository,
file_hash=file_hash,
)
return UploadAssetResponse(
storage_key=storage_key,
ingest_job_id=job.id,
asset_id=pending_asset.id,
url=file_url,
)