54916aff86
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (push) Successful in 5s
CI/CD Pipeline / Check push changed paths (push) Successful in 8s
CI/CD Pipeline / Build Staging Web Image (push) Successful in 29s
CI/CD Pipeline / Build Staging API Image (push) Successful in 31s
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 37s
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 2m1s
CI/CD Pipeline / Integration Tests (push) Successful in 2m22s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m11s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 1m16s
CI/CD Pipeline / Validate - Style (push) Successful in 2m58s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m41s
CI/CD Pipeline / ACR Image Cleanup (push) Successful in 1m30s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 1m36s
CI/CD Pipeline / Frontend Unit Tests (push) Successful in 5m47s
CI/CD Pipeline / Validate - Security (push) Successful in 6m2s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 3m33s
AI Code Review / AI Code Review (pull_request) Successful in 6m27s
CI/CD Pipeline / Unit Tests (push) Successful in 8m30s
CI/CD Pipeline / Production Browser E2E (push) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 2s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 31s
CI/CD Pipeline / PR Build Web Image (pull_request) Successful in 32s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 38s
CI/CD Pipeline / Production Browser E2E (pull_request) Has been cancelled
CI/CD Pipeline / CI Gate (pull_request) Has been cancelled
CI/CD Pipeline / ACR Image Cleanup (pull_request) Failing after 241h24m45s
CI/CD Pipeline / Staging API Integration Tests (pull_request) Failing after 241h24m46s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Failing after 241h24m56s
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Failing after 241h25m4s
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Failing after 241h25m7s
CI/CD Pipeline / Build Staging Worker Image (pull_request) Failing after 241h25m17s
CI/CD Pipeline / Build Staging API Image (pull_request) Failing after 241h25m19s
CI/CD Pipeline / Check push changed paths (pull_request) Failing after 241h25m57s
CI/CD Pipeline / Canary Release to Production (push) Failing after 241h34m12s
CI/CD Pipeline / CI Gate (push) Failing after 241h34m16s
CI/CD Pipeline / Build Production Worker Image (push) Failing after 241h34m16s
CI/CD Pipeline / Build Production API Image (push) Failing after 241h34m16s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Failing after 241h38m54s
CI/CD Pipeline / Canary Release to Production (pull_request) Failing after 0s
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Failing after 241h41m44s
CI/CD Pipeline / Retag skipped Staging API Image (push) Failing after 241h41m46s
CI/CD Pipeline / Build Production Worker Image (pull_request) Failing after 0s
CI/CD Pipeline / Build Production API Image (pull_request) Failing after 0s
CI/CD Pipeline / Frontend Unit Tests (pull_request) Failing after 241h25m51s
CI/CD Pipeline / Integration Tests (pull_request) Failing after 241h25m52s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 241h25m52s
CI/CD Pipeline / Validate - Security (pull_request) Failing after 241h25m56s
CI/CD Pipeline / Validate - Style (pull_request) Failing after 241h25m56s
CI/CD Pipeline / Frontend Lint (push) Failing after 241h42m46s
CI/CD Pipeline / PR Build Worker Image (push) Failing after 241h42m51s
CI/CD Pipeline / PR Build API Image (push) Failing after 241h42m52s
CI/CD Pipeline / Check if frontend-only change (push) Failing after 241h42m57s
CI/CD Pipeline / Staging E2E Tests (pull_request) Failing after 241h59m22s
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Failing after 241h59m42s
CI/CD Pipeline / Build Staging Web Image (pull_request) Failing after 241h59m53s
CI/CD Pipeline / Deploy Production (push) Failing after 242h8m49s
CI/CD Pipeline / Build Production Web Image (push) Failing after 242h8m52s
CI/CD Pipeline / Deploy Production (pull_request) Failing after 0s
CI/CD Pipeline / Retag skipped Staging Web Image (push) Failing after 242h16m21s
CI/CD Pipeline / Build Production Web Image (pull_request) Failing after 0s
CI/CD Pipeline / Frontend Lint (pull_request) Failing after 242h0m28s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Failing after 242h0m28s
CI/CD Pipeline / PR Build Web Image (push) Failing after 242h17m27s
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com> Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
587 lines
22 KiB
Python
587 lines
22 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
|
||
|
||
|
||
# 兜底去重:无 file_hash / client_upload_id 且大小已知时,同库同名同大小近期活动记录视为重复
|
||
FALLBACK_DEDUP_WINDOW_MINUTES = 30
|
||
|
||
|
||
def _find_duplicate_asset(
|
||
asset_repository: Any,
|
||
*,
|
||
library_id: str,
|
||
file_hash: str,
|
||
client_upload_id: str,
|
||
filename: str,
|
||
file_size: int = 0,
|
||
) -> Any:
|
||
"""complete/上传幂等去重,按优先级查找已存在的素材。
|
||
|
||
1. client_upload_id(客户端幂等 token,同一次上传的重试保持一致)
|
||
2. file_hash(内容哈希,不同上传只要内容相同即去重)
|
||
3. 兜底(严格模式,宁可漏判不可误杀):file_hash 与 client_upload_id
|
||
均缺失、且 file_size > 0 时,同库 + 同文件名 + **同大小** 且 30 分钟内
|
||
仍处 uploading/processing 的记录才判重。
|
||
- file_hash 非空时跳过兜底(hash 已代表内容;同名但内容全新的视频
|
||
如 iPhone 的 IMG_xxxx.MOV 绝不能被同名占位误杀)
|
||
- file_size=0(未知)时不允许仅凭同名 + processing 判重,直接放行
|
||
|
||
全部为鸭子类型调用:旧仓储无对应方法时静默跳过,不破坏既有实现。
|
||
"""
|
||
if client_upload_id:
|
||
find = getattr(asset_repository, "find_by_library_and_client_upload_id", None)
|
||
if callable(find):
|
||
existing = find(library_id=library_id, client_upload_id=client_upload_id)
|
||
if existing is not None:
|
||
logger.info(
|
||
"素材幂等命中(client_upload_id): library=%s token=%s asset=%s",
|
||
library_id,
|
||
client_upload_id,
|
||
getattr(existing, "id", "?"),
|
||
)
|
||
return existing
|
||
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(
|
||
"素材去重命中(file_hash): library=%s hash=%s asset=%s",
|
||
library_id,
|
||
file_hash,
|
||
existing.id,
|
||
)
|
||
return existing
|
||
# 同名兜底去重(最后防线,严格模式):
|
||
# - 仅当 file_hash / client_upload_id 均缺失时启用(hash 能代表内容时不靠同名猜)
|
||
# - file_size 必须 > 0 且与记录大小严格一致;大小未知(0)直接放行
|
||
# - 只命中近期 UPLOADING/PROCESSING 活动记录(READY 历史素材不拦)
|
||
if filename and not file_hash and not client_upload_id and file_size and file_size > 0:
|
||
find_recent = getattr(asset_repository, "find_recent_active_by_library_and_name", None)
|
||
if callable(find_recent):
|
||
existing = find_recent(
|
||
library_id=library_id,
|
||
name=filename,
|
||
within_minutes=FALLBACK_DEDUP_WINDOW_MINUTES,
|
||
file_size=file_size,
|
||
)
|
||
if existing is not None:
|
||
logger.info(
|
||
"素材幂等兜底命中(近期同名同大小活动记录): library=%s name=%s asset=%s status=%s size=%s",
|
||
library_id,
|
||
filename,
|
||
getattr(existing, "id", "?"),
|
||
getattr(existing, "status", None),
|
||
file_size,
|
||
)
|
||
return existing
|
||
elif filename and not file_hash and not client_upload_id and not file_size:
|
||
logger.debug(
|
||
"同名兜底去重跳过(file_size 未知,宁可放行不可误杀): library=%s name=%s",
|
||
library_id,
|
||
filename,
|
||
)
|
||
return None
|
||
|
||
|
||
def _create_pending_asset(
|
||
asset_repository,
|
||
project_id,
|
||
library_id,
|
||
storage_key,
|
||
filename,
|
||
mime_type,
|
||
user_id,
|
||
file_hash="",
|
||
client_upload_id="",
|
||
file_size: int = 0,
|
||
):
|
||
"""立即创建或复用一条 PROCESSING 状态的 Asset 记录。
|
||
|
||
find-or-create:prepare 阶段已按 file_hash/client_upload_id 预建的占位记录
|
||
会被 find_by_library_and_file_hash/find_by_library_and_client_upload_id 命中,
|
||
直接复用并补齐字段(避免 pre-create + complete 重复建两条)。
|
||
"""
|
||
# 1. 按 client_upload_id / file_hash 查找现有记录
|
||
existing = None
|
||
if client_upload_id:
|
||
find_by_cuid = getattr(asset_repository, "find_by_library_and_client_upload_id", None)
|
||
if callable(find_by_cuid):
|
||
existing = find_by_cuid(library_id=library_id, client_upload_id=client_upload_id)
|
||
if existing is None and file_hash:
|
||
existing = asset_repository.find_by_library_and_file_hash(library_id=library_id, file_hash=file_hash)
|
||
if existing is not None:
|
||
# 补齐字段(幂等:避免重复建记录,前端已拿到 asset_id)
|
||
changed = False
|
||
if file_hash and not existing.file_hash:
|
||
existing.file_hash = file_hash
|
||
changed = True
|
||
if client_upload_id and not existing.client_upload_id:
|
||
existing.client_upload_id = client_upload_id
|
||
changed = True
|
||
if file_size and not existing.file_size:
|
||
existing.file_size = file_size
|
||
changed = True
|
||
if existing.status not in (AssetStatus.PROCESSING, AssetStatus.UPLOADING):
|
||
existing.status = AssetStatus.PROCESSING
|
||
changed = True
|
||
if changed:
|
||
try:
|
||
asset_repository.update(existing)
|
||
except Exception: # noqa: BLE001 — 字段补齐失败不阻塞主流程
|
||
pass
|
||
return existing
|
||
|
||
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,
|
||
client_upload_id=client_upload_id,
|
||
file_size=file_size,
|
||
)
|
||
return asset_repository.create(asset)
|
||
|
||
|
||
def _persist_celery_task_id(repo: Any, job: Any, celery_task_id: str) -> None:
|
||
"""记录 celery 消息 ID 到任务行,供孤儿清理时 revoke/清除队列消息(#1714)。"""
|
||
if not celery_task_id:
|
||
return
|
||
try:
|
||
job.celery_task_id = celery_task_id
|
||
repo.update(job)
|
||
except Exception: # noqa: BLE001 — 记录失败不影响主流程(执行前状态守卫兜底)
|
||
pass
|
||
|
||
|
||
def _submit_ingest_job(
|
||
project_id: str,
|
||
library_id: str,
|
||
storage_key: str,
|
||
ingest_job_repository: Any,
|
||
file_hash: str = "",
|
||
asset_id: 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,
|
||
asset_id=asset_id,
|
||
)
|
||
)
|
||
celery_result = celery_app.send_task("worker.ingest_asset", args=[job.id])
|
||
_persist_celery_task_id(ingest_job_repository, job, getattr(celery_result, "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),
|
||
asset_repository: Any = Depends(get_asset_repository),
|
||
storage_service: OSSStorageService = Depends(get_storage_service),
|
||
) -> DirectUploadPrepareResponse:
|
||
"""创建浏览器直传 OSS 的短期表单签名,并在签名前按 file_hash/client_upload_id 去重。
|
||
|
||
命中去重:直接返回 duplicated=True + skip_transfer=True(前端跳过 OSS 直传),
|
||
未命中:正常签名 OSS 并立即预建一条 PROCESSING 状态的 asset 记录占住
|
||
file_hash 闸门,响应带 asset_id 供前端/后续 complete 关联。
|
||
"""
|
||
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,
|
||
)
|
||
|
||
safe_filename = request.filename.replace("/", "_").replace("\\", "_")
|
||
|
||
# ── prepare 阶段去重:OSS 签名之前先查已存在素材 ──
|
||
if request.file_hash or request.client_upload_id:
|
||
existing = _find_duplicate_asset(
|
||
asset_repository,
|
||
library_id=request.library_id,
|
||
file_hash=request.file_hash,
|
||
client_upload_id=request.client_upload_id,
|
||
filename=request.filename,
|
||
file_size=request.file_size,
|
||
)
|
||
if existing is not None:
|
||
logger.info(
|
||
"prepare 命中去重: library=%s hash=%s cuid=%s existing_asset=%s",
|
||
request.library_id,
|
||
request.file_hash,
|
||
request.client_upload_id,
|
||
existing.id,
|
||
)
|
||
return DirectUploadPrepareResponse(
|
||
upload_url="",
|
||
method="",
|
||
storage_key=existing.storage_key,
|
||
expires_at="",
|
||
fields={},
|
||
max_size_bytes=0,
|
||
duplicated=True,
|
||
skip_transfer=True,
|
||
asset_id=existing.id,
|
||
)
|
||
|
||
file_id = uuid4().hex[:8]
|
||
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
|
||
|
||
# ── 预建 asset 占位:占住 file_hash/client_upload_id 闸门,避免并发重复上传 ──
|
||
pending_asset_id = ""
|
||
if request.file_hash or request.client_upload_id:
|
||
try:
|
||
pending = _create_pending_asset(
|
||
asset_repository=asset_repository,
|
||
project_id=request.project_id,
|
||
library_id=request.library_id,
|
||
storage_key=storage_key,
|
||
filename=safe_filename,
|
||
mime_type=validated_content_type,
|
||
user_id=authenticated_user.user.id,
|
||
file_hash=request.file_hash,
|
||
client_upload_id=request.client_upload_id,
|
||
file_size=request.file_size,
|
||
)
|
||
pending_asset_id = pending.id
|
||
except Exception as error:
|
||
# 预建失败不阻塞签名:complete 仍可按 OSS 文件 + hash 兜底去重
|
||
logger.warning("预建 asset 占位失败,降级走 old flow: %s", 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,
|
||
duplicated=False,
|
||
skip_transfer=False,
|
||
asset_id=pending_asset_id,
|
||
)
|
||
|
||
|
||
@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:
|
||
"""确认浏览器直传完成并创建导入任务(幂等:重复 complete 返回同一素材)。"""
|
||
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")
|
||
|
||
filename = normalized_key.rsplit("/", 1)[-1]
|
||
|
||
# ── 幂等去重(放在 OSS 检查之前):complete 超时后前端重试时,
|
||
# 第一次 complete 可能已建好占位记录,此时即使 OSS 检查失败也必须返回
|
||
# 已存在记录,绝不能再建第二条。─
|
||
existing = _find_duplicate_asset(
|
||
asset_repository,
|
||
library_id=request.library_id,
|
||
file_hash=request.file_hash,
|
||
client_upload_id=request.client_upload_id,
|
||
filename=filename,
|
||
file_size=request.file_size,
|
||
)
|
||
if existing is not None:
|
||
return DirectUploadCompleteResponse(
|
||
storage_key=existing.storage_key,
|
||
ingest_job_id="",
|
||
duplicated=True,
|
||
asset_id=existing.id,
|
||
url=storage_service.get_url(existing.storage_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")
|
||
|
||
# 立即创建 Asset 记录(PROCESSING 状态),使前端刷新后即可看到新素材
|
||
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,
|
||
client_upload_id=request.client_upload_id,
|
||
file_size=request.file_size,
|
||
)
|
||
|
||
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,
|
||
asset_id=pending_asset.id,
|
||
)
|
||
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="文件哈希,用于去重检测"),
|
||
client_upload_id: str = Form(default="", description="客户端幂等 token(同一次上传的重试保持一致)"),
|
||
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)
|
||
|
||
# P2-5: 服务端验证 MIME 类型(先验证,再幂等去重,避免非法类型绕过)
|
||
validated_content_type = _validate_mime_type(file.content_type)
|
||
|
||
safe_filename = file.filename.replace("/", "_").replace("\\", "_") if file.filename else "unknown"
|
||
|
||
# ── 幂等去重:client_upload_id → file_hash → 近期活动同名记录兜底 ──
|
||
# 放在 OSS 上传之前:重复提交直接返回,不占 OSS 流量、不建新记录。
|
||
existing = _find_duplicate_asset(
|
||
asset_repository,
|
||
library_id=library_id,
|
||
file_hash=file_hash,
|
||
client_upload_id=client_upload_id,
|
||
filename=safe_filename,
|
||
file_size=0,
|
||
)
|
||
if existing is not None:
|
||
return UploadAssetResponse(
|
||
storage_key=existing.storage_key,
|
||
ingest_job_id="",
|
||
url="",
|
||
duplicated=True,
|
||
asset_id=existing.id,
|
||
)
|
||
|
||
file_id = uuid4().hex[:8]
|
||
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,
|
||
client_upload_id=client_upload_id,
|
||
)
|
||
|
||
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,
|
||
asset_id=pending_asset.id,
|
||
)
|
||
|
||
return UploadAssetResponse(
|
||
storage_key=storage_key,
|
||
ingest_job_id=job.id,
|
||
asset_id=pending_asset.id,
|
||
url=file_url,
|
||
)
|