This commit is contained in:
2026-07-09 09:10:49 +08:00
parent 2882d37112
commit fec47bf47a
54 changed files with 5985 additions and 992 deletions
@@ -50,6 +50,7 @@ async def create_menu_config(
ensure_ascii=False,
),
)
await db.commit()
return menu
@@ -89,6 +90,7 @@ async def update_menu_config(
ensure_ascii=False,
),
)
await db.commit()
return menu
@@ -121,4 +123,5 @@ async def delete_menu_config(
ensure_ascii=False,
),
)
await db.commit()
return {"message": "ok"}
@@ -72,6 +72,7 @@ async def create_package(
ensure_ascii=False,
),
)
await db.commit()
return _to_out(pkg)
@@ -109,6 +110,7 @@ async def update_package(
ensure_ascii=False,
),
)
await db.commit()
return _to_out(pkg)
@@ -142,4 +144,5 @@ async def delete_package(
ensure_ascii=False,
),
)
await db.commit()
return {"ok": True}
@@ -58,6 +58,7 @@ async def save_config(
ensure_ascii=False,
),
)
await db.commit()
return result
@@ -81,6 +82,7 @@ async def reset_default(
ensure_ascii=False,
),
)
await db.commit()
return result
@@ -114,6 +116,7 @@ async def import_config(
ensure_ascii=False,
),
)
await db.commit()
return result
+2
View File
@@ -34,6 +34,7 @@ from app.api.admin import router as admin_module_router
from app.api.v1.material_admin import router as material_admin_router
from app.api.v1.private_portrait import router as private_portrait_router
from app.api.v1.private_portrait_virtual import router as private_portrait_virtual_router
from app.api.v1.upload_resource import router as upload_resource_router
api_router = APIRouter()
api_router.include_router(auth_router)
@@ -70,3 +71,4 @@ api_router.include_router(admin_module_router)
api_router.include_router(material_admin_router)
api_router.include_router(private_portrait_router)
api_router.include_router(private_portrait_virtual_router)
api_router.include_router(upload_resource_router)
+10 -3
View File
@@ -202,6 +202,7 @@ async def create_user(
),
ip=None,
)
await db.commit()
return user
@@ -233,6 +234,7 @@ async def update_user_menus(
ensure_ascii=False,
),
)
await db.commit()
return {"message": "ok"}
@@ -289,6 +291,7 @@ async def adjust_credits(
ensure_ascii=False,
),
)
await db.commit()
return {"message": "ok"}
@@ -319,6 +322,7 @@ async def update_user_status(
ensure_ascii=False,
),
)
await db.commit()
return {"message": "ok"}
@@ -1594,7 +1598,7 @@ async def update_system_config(
if not config:
raise HTTPException(status_code=404, detail="配置不存在")
config.value = str(req.value)
await db.commit()
await db.flush()
await log_operation(
db,
admin.id,
@@ -1611,6 +1615,7 @@ async def update_system_config(
ensure_ascii=False,
),
)
await db.commit()
return config
@@ -2185,7 +2190,7 @@ async def upload_pdf(
key=config_key,
value=url,
))
await db.commit()
await db.flush()
await log_operation(
db,
admin.id,
@@ -2202,6 +2207,7 @@ async def upload_pdf(
ensure_ascii=False,
),
)
await db.commit()
return {"url": url}
@@ -2247,7 +2253,7 @@ async def upload_logo(
value=url,
description="网站Logo图片",
))
await db.commit()
await db.flush()
await log_operation(
db,
admin.id,
@@ -2263,6 +2269,7 @@ async def upload_logo(
ensure_ascii=False,
),
)
await db.commit()
return {"url": url}
+155 -92
View File
@@ -36,6 +36,9 @@ from app.services.resource_accounting_service import (
from app.services.private_portrait.reference_resolver import batch_resolve_private_portrait_reference_display_urls, resolve_private_portrait_reference_display_urls
from app.services.resource_signed_url_service import build_resource_signed_url
from app.services.resource_capacity_service import assert_user_resource_capacity_available
from app.services.upload_resource import delete_unbound_upload_resource, upload_reference_file, cleanup_upload_resource_files_after_commit
from app.services.upload_resource.log_service import log_upload_resource_exception, safe_rollback_with_log
from app.enums.upload_resource import UploadResourceEventEnum, UploadResourceModuleEnum, UploadResourceTypeEnum
from app.services.generation_billing_service import (
CHARGE_TEXT_PROMPT,
OWNER_GENERATION_RECORD,
@@ -780,128 +783,188 @@ async def seedance_callback(request: Request, db: AsyncSession = Depends(get_db)
return {"message": "ok"}
@router.post("/upload-image")
@router.post(
"/upload-image",
summary="上传 AI 创作普通参考图片",
description=(
"上传普通 AI 创作参考图片,写入 UploadResource 资源账本并纳入用户上传容量统计。"
"返回 resource_id 和 url。该文件在未绑定业务记录前可通过 /generation-records/delete-file 单独删除,"
"也会出现在 /upload-resources/history 历史素材中供 AI 创作复用。"
),
responses={400: {"description": "文件类型、大小或容量校验失败"}, 401: {"description": "未登录或 Token 无效"}},
)
async def upload_image(
file: UploadFile = File(...),
current_user: User = Depends(get_current_user),
gen_type: str = Query("video", description="生成类型:video-视频,image-图片"),
db: AsyncSession = Depends(get_db),
):
"""Upload an image for generation reference."""
import os
import uuid
from app.config import settings
from datetime import datetime
if not file.content_type or not file.content_type.startswith("image/"):
raise HTTPException(status_code=400, detail="仅支持图片文件")
ext = os.path.splitext(file.filename or ".png")[1] or ".png"
timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
safe_name = f"{gen_type}_img_{current_user.id}_{timestamp}_{uuid.uuid4().hex[:8]}{ext}"
date_dir = datetime.now().strftime("%Y/%m/%d")
dir_path = os.path.join(settings.UPLOAD_LOCAL_PATH, "images", date_dir)
os.makedirs(dir_path, exist_ok=True)
file_path = os.path.join(dir_path, safe_name)
content = await file.read()
if len(content) > 10 * 1024 * 1024:
raise HTTPException(status_code=400, detail="图片大小不能超过10MB")
with open(file_path, "wb") as f:
f.write(content)
url = f"/uploads/images/{date_dir}/{safe_name}"
return {"url": url, "filename": file.filename or safe_name, "type": "image", "gen_type": gen_type}
"""上传普通参考图片,记录 UploadResource 并纳入用户容量统计。"""
result = await upload_reference_file(
db,
file=file,
current_user=current_user,
module=UploadResourceModuleEnum.COMMON.value,
resource_type=UploadResourceTypeEnum.IMAGE.value,
gen_type=gen_type,
)
await db.commit()
return {
"url": result.url,
"filename": result.filename,
"type": "image",
"gen_type": gen_type,
"resource_id": result.resource_id,
"file_size_bytes": result.file_size_bytes,
}
@router.post("/upload-video")
@router.post(
"/upload-video",
summary="上传 AI 创作普通参考视频",
description=(
"上传普通 AI 创作参考视频,写入 UploadResource 资源账本并纳入用户上传容量统计。"
"duration_seconds 为前端识别的视频秒数,用于历史复用和 AI 创作视频总时长校验。"
"未绑定业务记录前可单独删除,并会出现在 /upload-resources/history 历史素材中。"
),
responses={400: {"description": "文件类型、大小、容量或视频参数校验失败"}, 401: {"description": "未登录或 Token 无效"}},
)
async def upload_video(
file: UploadFile = File(...),
current_user: User = Depends(get_current_user),
duration_seconds: float | None = Query(None, description="前端识别的视频时长秒数,可选"),
db: AsyncSession = Depends(get_db),
):
"""Upload a video for generation reference."""
import os
import uuid
from app.config import settings
from datetime import datetime
if not file.content_type or not file.content_type.startswith("video/"):
raise HTTPException(status_code=400, detail="仅支持视频文件")
ext = os.path.splitext(file.filename or ".mp4")[1] or ".mp4"
timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
safe_name = f"video_ref_{current_user.id}_{timestamp}_{uuid.uuid4().hex[:8]}{ext}"
date_dir = datetime.now().strftime("%Y/%m/%d")
dir_path = os.path.join(settings.UPLOAD_LOCAL_PATH, "videos", date_dir)
os.makedirs(dir_path, exist_ok=True)
file_path = os.path.join(dir_path, safe_name)
content = await file.read()
if len(content) > 100 * 1024 * 1024:
raise HTTPException(status_code=400, detail="视频大小不能超过100MB")
with open(file_path, "wb") as f:
f.write(content)
url = f"/uploads/videos/{date_dir}/{safe_name}"
return {"url": url, "filename": file.filename or safe_name, "type": "video"}
"""上传普通参考视频,记录 UploadResource 并纳入用户容量统计。"""
result = await upload_reference_file(
db,
file=file,
current_user=current_user,
module=UploadResourceModuleEnum.COMMON.value,
resource_type=UploadResourceTypeEnum.VIDEO.value,
duration_seconds=duration_seconds,
)
await db.commit()
return {
"url": result.url,
"filename": result.filename,
"type": "video",
"resource_id": result.resource_id,
"file_size_bytes": result.file_size_bytes,
"duration_seconds": result.duration_seconds,
}
@router.post("/upload-audio")
@router.post(
"/upload-audio",
summary="上传 AI 创作普通参考音频",
description=(
"上传普通 AI 创作参考音频,写入 UploadResource 资源账本并纳入用户上传容量统计。"
"当前仅支持 mp3、wav;单文件大小受 AUDIO_MAX_FILE_SIZE_MB 限制。"
"duration_seconds 为前端识别的音频秒数,用于 AI 创作音频总时长校验。"
"未绑定业务记录前可单独删除,并会出现在 /upload-resources/history 历史素材中。"
),
responses={400: {"description": "音频格式、MIME、大小或容量校验失败"}, 401: {"description": "未登录或 Token 无效"}},
)
async def upload_audio(
file: UploadFile = File(...),
current_user: User = Depends(get_current_user),
duration_seconds: float | None = Query(None, description="前端识别的音频时长秒数,可选"),
db: AsyncSession = Depends(get_db),
):
"""Upload an audio file for AI creation reference."""
"""上传普通参考音频,记录 UploadResource 并纳入用户容量统计。"""
import os
import uuid
from app.config import settings
from datetime import datetime
ext = os.path.splitext(file.filename or "")[1].lower().lstrip(".")
if ext not in AUDIO_ALLOWED_EXTENSIONS:
raise HTTPException(status_code=400, detail="仅支持 mp3、wav 音频文件")
expected_mime = AUDIO_ALLOWED_MIME_TYPES.get(ext)
if not file.content_type or file.content_type != expected_mime:
if expected_mime and file.content_type and file.content_type != expected_mime:
raise HTTPException(status_code=400, detail=f"音频 MIME 类型错误,{ext} 必须为 {expected_mime}")
timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
safe_name = f"audio_ref_{current_user.id}_{timestamp}_{uuid.uuid4().hex[:8]}.{ext}"
date_dir = datetime.now().strftime("%Y/%m/%d")
dir_path = os.path.join(settings.UPLOAD_LOCAL_PATH, "audios", date_dir)
os.makedirs(dir_path, exist_ok=True)
file_path = os.path.join(dir_path, safe_name)
content = await file.read()
max_bytes = AUDIO_MAX_FILE_SIZE_MB * 1024 * 1024
if len(content) > max_bytes:
raise HTTPException(status_code=400, detail=f"音频大小不能超过{AUDIO_MAX_FILE_SIZE_MB}MB")
with open(file_path, "wb") as f:
f.write(content)
url = f"/uploads/audios/{date_dir}/{safe_name}"
return {"url": url, "filename": file.filename or safe_name, "type": "audio"}
result = await upload_reference_file(
db,
file=file,
current_user=current_user,
module=UploadResourceModuleEnum.COMMON.value,
resource_type=UploadResourceTypeEnum.AUDIO.value,
duration_seconds=duration_seconds,
max_bytes=AUDIO_MAX_FILE_SIZE_MB * 1024 * 1024,
)
await db.commit()
return {
"url": result.url,
"filename": result.filename,
"type": "audio",
"resource_id": result.resource_id,
"file_size_bytes": result.file_size_bytes,
"duration_seconds": result.duration_seconds,
}
@router.post("/delete-file")
@router.post(
"/delete-file",
summary="删除未绑定上传文件",
description=(
"删除当前用户自己的未绑定上传文件,并释放 UploadResource 上传容量。"
"仅允许删除 bind_status=pending、delete_policy=user_deletable、未绑定 source_model/source_id 的资源。"
"删除顺序为主事务先 soft delete 并 commitcommit 成功后再清理真实文件。"
"该接口保留给单文件删除;批量删除请使用 DELETE /upload-resources/history/batch。"
),
responses={400: {"description": "文件路径无效、文件已被模块任务使用或不可单独删除"}, 401: {"description": "未登录或 Token 无效"}, 403: {"description": "无权删除此文件"}},
)
async def delete_upload(
url: str = Query(..., description="文件URL,如 /uploads/images/2024/01/01/video_img_xxx.png"),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
"""Delete an uploaded file by URL."""
import os
from app.config import settings
"""删除未绑定业务记录的上传文件,并释放 UploadResource 容量。"""
user_id = str(current_user.id)
pending_ids: list[str] = []
legacy_paths: list[str] = []
try:
result = await delete_unbound_upload_resource(db, user=current_user, url=url)
pending_ids = list(result.pop("_pending_physical_delete_resource_ids", []) or [])
legacy_paths = list(result.pop("_legacy_pending_delete_paths", []) or [])
await db.commit()
except Exception as exc: # noqa: BLE001
await safe_rollback_with_log(
db,
event_type=UploadResourceEventEnum.DELETE_UPLOAD_ROLLBACK_FAILED.value,
message="删除上传文件主事务回滚失败",
user_id=user_id,
detail={"url": url},
original_exc=exc,
)
log_upload_resource_exception(
event_type=UploadResourceEventEnum.DELETE_UPLOAD_FAILED.value,
message=f"删除上传文件失败: {exc}",
user_id=user_id,
detail={"url": url},
exc=exc,
)
raise
if not url.startswith("/uploads/"):
raise HTTPException(status_code=400, detail="无效的文件路径")
if current_user.id not in url:
raise HTTPException(status_code=403, detail="无权删除此文件")
file_path = os.path.join(settings.UPLOAD_LOCAL_PATH, url.replace("/uploads/", ""))
if os.path.exists(file_path):
os.remove(file_path)
return {"message": "ok"}
if pending_ids or legacy_paths:
try:
await cleanup_upload_resource_files_after_commit(db, resource_ids=pending_ids, legacy_paths=legacy_paths)
await db.commit()
except Exception as exc: # noqa: BLE001
await safe_rollback_with_log(
db,
event_type=UploadResourceEventEnum.DELETE_UPLOAD_ROLLBACK_FAILED.value,
message="删除上传文件 cleanup 事务回滚失败",
user_id=user_id,
detail={"url": url, "pending_ids": pending_ids, "legacy_paths": legacy_paths},
original_exc=exc,
)
log_upload_resource_exception(
event_type=UploadResourceEventEnum.DELETE_UPLOAD_CLEANUP_FAILED.value,
message=f"删除上传文件后清理真实文件失败: {exc}",
user_id=user_id,
resource_ids=pending_ids,
detail={"url": url, "legacy_paths": legacy_paths},
exc=exc,
)
return result
+11 -2
View File
@@ -121,6 +121,9 @@ async def list_engines(
"该接口用于 Chat 风格的图片/视频生成,不再绑定 project_id。"
"创建成功后会写入 chat_generation_tasks 表,并投递 Celery 异步任务。"
"支持 image 图片生成和 video 视频生成。"
"media_references 支持 image/video/audiosource 可为 upload_resource=历史上传素材、"
"private_portrait_asset=私域真人/虚拟素材、空=本次普通上传素材。"
"视频/音频参考素材必须携带 duration,并按 AI 创作原逻辑校验单段 2~15 秒、总时长不超过 15 秒。"
"建议前端传入 idempotency_key,用于防止按钮连点、网络重试导致重复创建任务和重复扣费。"
),
responses={
@@ -141,7 +144,12 @@ async def list_engines(
async def create_task(
req: GenerationAITaskCreate = Body(
...,
description="AI生成任务创建参数。gen_type=image 时使用图片参数;gen_type=video 时使用视频参数",
description=(
"AI生成任务创建参数。gen_type=image 时使用图片参数;gen_type=video 时使用视频参数。"
"枚举:gen_type=image/videomedia_references[].type=image/video/audio"
"media_references[].source=upload_resource/private_portrait_asset/空;"
"media_references[].role=first_frame/last_frame/reference_image/reference_video/reference_audio。"
),
),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
@@ -186,7 +194,8 @@ async def create_task(
"分页获取当前登录用户的AI生成任务列表。"
"可按生成类型 gen_type 和任务状态 status 过滤。"
"该接口返回的是普通任务列表,不按日期分组。"
"如果前端需要按生成日期分组展示历史记录,请使用 /generation-ai/history 接口。"
"如果前端需要按生成日期分组展示生成历史记录,请使用 /generation-ai/history 接口。"
"上传素材历史不是生成历史,请使用 /upload-resources/history。"
),
responses={
200: {
@@ -3,7 +3,7 @@ from __future__ import annotations
from datetime import datetime
from types import SimpleNamespace
from fastapi import APIRouter, Body, Depends, HTTPException, Path, Query
from fastapi import APIRouter, Body, Depends, File, HTTPException, Path, Query, UploadFile
from sqlalchemy import inspect as sa_inspect
from sqlalchemy.ext.asyncio import AsyncSession
@@ -48,6 +48,9 @@ from app.services.module_async_recovery_service import (
register_module_step_task,
)
from app.tasks.celery_app import celery_app
from app.enums.upload_resource import UploadResourceEventEnum, UploadResourceModuleEnum, UploadResourceSourceModelEnum, UploadResourceTypeEnum
from app.services.upload_resource import upload_reference_file, bind_upload_resources, cleanup_upload_resource_files_after_commit
from app.services.upload_resource.log_service import log_upload_resource_exception, safe_rollback_with_log
MODULE = ModuleCodeEnum.HOT_OPENING_REPLICATE.value
@@ -203,6 +206,66 @@ async def get_spec():
return HotOpeningSpecOut()
@router.post(
"/upload-image",
summary="上传爆款开头复刻图片素材",
description="上传后先记录为 pending UploadResource;创建项目成功后绑定 ModuleGenerationProject。",
)
async def upload_hot_opening_image(
file: UploadFile = File(...),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
result = await upload_reference_file(
db,
file=file,
current_user=current_user,
module=UploadResourceModuleEnum.HOT_OPENING_REPLICATE.value,
resource_type=UploadResourceTypeEnum.IMAGE.value,
gen_type="video",
)
await db.commit()
return {
"url": result.url,
"filename": result.filename,
"type": "image",
"module": result.module,
"resource_id": result.resource_id,
"file_size_bytes": result.file_size_bytes,
}
@router.post(
"/upload-video",
summary="上传爆款开头复刻视频素材",
description="上传后先记录为 pending UploadResource;创建项目成功后绑定 ModuleGenerationProject。",
)
async def upload_hot_opening_video(
file: UploadFile = File(...),
duration_seconds: float | None = Query(None, description="前端识别的视频时长秒数,可选"),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
result = await upload_reference_file(
db,
file=file,
current_user=current_user,
module=UploadResourceModuleEnum.HOT_OPENING_REPLICATE.value,
resource_type=UploadResourceTypeEnum.VIDEO.value,
duration_seconds=duration_seconds,
)
await db.commit()
return {
"url": result.url,
"filename": result.filename,
"type": "video",
"module": result.module,
"resource_id": result.resource_id,
"file_size_bytes": result.file_size_bytes,
"duration_seconds": result.duration_seconds,
}
@router.post(
"/tasks",
response_model=HotOpeningTaskDetailOut,
@@ -222,6 +285,16 @@ async def create_task(
try:
project = await create_hot_opening_project(db, current_user, req)
project_id_value = str(project.id)
await bind_upload_resources(
db,
user_id=current_user.id,
module=UploadResourceModuleEnum.HOT_OPENING_REPLICATE.value,
source_model=UploadResourceSourceModelEnum.MODULE_GENERATION_PROJECT.value,
source_id=project_id_value,
resource_ids=[req.material_video_resource_id, req.material_image_resource_id],
urls=[req.material_video_url, req.material_image_url],
allow_common_migrate=True,
)
await db.commit()
except HTTPException:
await db.rollback()
@@ -703,13 +776,72 @@ async def delete_task(
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
user_id = _safe_user_id(current_user)
pending_ids: list[str] = []
try:
result = await delete_hot_opening_project(db, current_user=current_user, project_id=project_id)
pending_ids = list(result.pending_delete_resource_ids or [])
await db.commit()
return result
except HTTPException:
await db.rollback()
except HTTPException as exc:
await safe_rollback_with_log(
db,
event_type=UploadResourceEventEnum.HOT_OPENING_DELETE_PROJECT_ROLLBACK_FAILED.value,
message="删除爆款开头复刻项目 HTTPException 回滚失败",
user_id=user_id,
module=MODULE,
detail={"project_id": project_id},
original_exc=exc,
)
raise
except Exception as exc:
await db.rollback()
except Exception as exc: # noqa: BLE001
await safe_rollback_with_log(
db,
event_type=UploadResourceEventEnum.HOT_OPENING_DELETE_PROJECT_ROLLBACK_FAILED.value,
message="删除爆款开头复刻项目主事务回滚失败",
user_id=user_id,
module=MODULE,
detail={"project_id": project_id},
original_exc=exc,
)
_log_api_error(
event_type=HotOpeningLogEventEnum.API_REQUEST_FAILED.value,
current_user=current_user,
project_id=project_id,
message=f"删除爆款开头复刻项目失败: {exc}",
detail={"project_id": project_id},
exc=exc,
)
log_upload_resource_exception(
event_type=UploadResourceEventEnum.HOT_OPENING_DELETE_PROJECT_FAILED.value,
message=f"删除爆款开头复刻项目失败: {exc}",
user_id=user_id,
module=MODULE,
detail={"project_id": project_id},
exc=exc,
)
raise HTTPException(status_code=500, detail=f"删除爆款开头复刻项目失败: {exc}")
if pending_ids:
try:
await cleanup_upload_resource_files_after_commit(db, resource_ids=pending_ids)
await db.commit()
except Exception as exc: # noqa: BLE001
await safe_rollback_with_log(
db,
event_type=UploadResourceEventEnum.HOT_OPENING_DELETE_PROJECT_ROLLBACK_FAILED.value,
message="删除爆款开头复刻项目 cleanup 事务回滚失败",
user_id=user_id,
module=MODULE,
detail={"project_id": project_id, "pending_ids": pending_ids},
original_exc=exc,
)
log_upload_resource_exception(
event_type=UploadResourceEventEnum.HOT_OPENING_DELETE_PROJECT_CLEANUP_FAILED.value,
message=f"删除爆款开头复刻项目后清理真实文件失败: {exc}",
user_id=user_id,
resource_ids=pending_ids,
module=MODULE,
detail={"project_id": project_id},
exc=exc,
)
return result
+143 -29
View File
@@ -9,7 +9,6 @@ from sqlalchemy.ext.asyncio import AsyncSession
from app.dependencies import get_current_user, get_db
from app.enums.private_portrait import (
PrivatePortraitAssetType,
PrivatePortraitEventSource,
PrivatePortraitEventStatus,
PrivatePortraitEventType,
@@ -23,8 +22,9 @@ from app.schemas.private_portrait import (
PrivatePortraitAssetCreate,
PrivatePortraitAssetListOut,
PrivatePortraitAssetOut,
PrivatePortraitDeleteOut,
PrivatePortraitConfigOut,
PrivatePortraitDeleteOut,
PrivatePortraitEnumMetaOut,
PrivatePortraitProjectCreate,
PrivatePortraitProjectCreateWithValidateOut,
PrivatePortraitProjectListOut,
@@ -33,7 +33,6 @@ from app.schemas.private_portrait import (
PrivatePortraitSelectableAssetListOut,
PrivatePortraitValidateSessionCreate,
PrivatePortraitValidateSessionOut,
PrivatePortraitEnumMetaOut,
build_private_portrait_enum_meta,
)
from app.services.operation_log_service import log_operation_error, log_operation_event
@@ -65,20 +64,47 @@ from app.services.private_portrait.real_person.service import (
router = APIRouter(tags=["私域真人素材库"])
_REAL_PERSON_API_DESCRIPTION = (
"私域真人素材库 API。真人和虚拟素材共用用户素材额度;user.private_portrait_asset_limit=0 表示关闭,"
">0 表示启用并限制总素材数量。真人项目创建后需要通过火山 CreateVisualValidateSession 进行人脸认证,"
"每个真人项目只允许认证一次。认证回调 resultCode=10000 表示成功,成功后才允许上传可用于 AI 创作的真人素材。"
)
def _log_task_dispatch_failed(*, task_name: str, user_id: str | None = None, project_id: str | None = None, asset_id: str | None = None, exc: BaseException) -> None:
log_operation_error(domain=DOMAIN, event_type=PrivatePortraitEventType.TASK_DISPATCH_FAILED.value, source=PrivatePortraitEventSource.API.value, user_id=user_id, project_id=project_id, asset_id=asset_id, exc=exc, detail={"task_name": task_name})
log_operation_error(
domain=DOMAIN,
event_type=PrivatePortraitEventType.TASK_DISPATCH_FAILED.value,
source=PrivatePortraitEventSource.API.value,
user_id=user_id,
project_id=project_id,
asset_id=asset_id,
exc=exc,
detail={"task_name": task_name},
)
def _log_task_dispatch_success(*, task_name: str, user_id: str | None = None, project_id: str | None = None, asset_id: str | None = None) -> None:
log_operation_event(domain=DOMAIN, event_type=PrivatePortraitEventType.TASK_DISPATCH_SUCCESS.value, event_status=PrivatePortraitEventStatus.SUCCESS.value, source=PrivatePortraitEventSource.API.value, user_id=user_id, project_id=project_id, asset_id=asset_id, detail={"task_name": task_name})
log_operation_event(
domain=DOMAIN,
event_type=PrivatePortraitEventType.TASK_DISPATCH_SUCCESS.value,
event_status=PrivatePortraitEventStatus.SUCCESS.value,
source=PrivatePortraitEventSource.API.value,
user_id=user_id,
project_id=project_id,
asset_id=asset_id,
detail={"task_name": task_name},
)
@router.get(
"/private-portrait/config",
response_model=PrivatePortraitConfigOut,
summary="获取当前用户私域人素材额度配置",
description="返回私域人像素材总量限制。额度由真人认证素材库与虚拟人像素材库共用,图片和视频共用,Audio 暂未开放。",
summary="获取当前用户私域人素材额度配置",
description=(
_REAL_PERSON_API_DESCRIPTION
+ "返回 enabled、asset_limit、used_asset_count、remaining_asset_count 等字段。额度按用户维度限制,真人/虚拟、图片/视频共用;Audio 暂不开放。"
),
)
async def get_my_private_portrait_config(current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
return await get_user_private_portrait_config(db, user_id=current_user.id)
@@ -87,8 +113,12 @@ async def get_my_private_portrait_config(current_user: User = Depends(get_curren
@router.get(
"/private-portrait/meta/enums",
response_model=PrivatePortraitEnumMetaOut,
summary="获取私域人素材库枚举说明",
description="给前端展示状态、类型、素材库类型使用。Audio 仅作为火山支持项展示,当前业务不开放上传。",
summary="获取私域人素材库枚举说明",
description=(
"返回前端展示所需枚举:library_type=real_person/aigc_virtual"
"asset_type=Image/Video/Audioproject status、asset status、remote_delete_status 等。"
"当前业务上传只开放 Image 和 Video,Audio 仅作为兼容枚举展示。"
),
)
async def get_private_portrait_enum_meta():
return build_private_portrait_enum_meta()
@@ -98,7 +128,11 @@ async def get_private_portrait_enum_meta():
"/private-portrait/projects",
response_model=PrivatePortraitProjectCreateWithValidateOut,
summary="创建真人认证素材项目并生成认证会话",
description="创建本地真人素材项目,随后调用火山 CreateVisualValidateSession 返回 H5Link。用户完成认证后,回调会创建本地 Asset Group 映射。",
description=(
_REAL_PERSON_API_DESCRIPTION
+ "创建本地真人项目后立即调用火山 CreateVisualValidateSession,返回 H5Link/BytedToken 给 PC 前端展示二维码。"
"手机扫码完成人脸认证后,PC 端通过 validate_session 查询状态;认证成功后才允许创建真人素材。"
),
)
async def create_private_portrait_project(payload: PrivatePortraitProjectCreate, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
project = await create_real_person_project(db, user_id=current_user.id, payload=payload)
@@ -112,13 +146,13 @@ async def create_private_portrait_project(payload: PrivatePortraitProjectCreate,
"/private-portrait/projects",
response_model=PrivatePortraitProjectListOut,
summary="查询当前用户真人认证素材项目列表",
description="只返回 library_type=real_person 的项目。默认查询 active 项目,可 status 覆盖。",
description="只返回 library_type=real_person 的项目。默认查询 active 项目,可通过 status 覆盖。",
)
async def list_private_portrait_projects(
page: int = Query(1, ge=1, description="页码,从 1 开始"),
page_size: int = Query(20, ge=1, le=100, description="每页数量,最大 100"),
keyword: str | None = Query(None, description="项目名称模糊搜索"),
status: str | None = Query(None, description="项目状态,不传默认 active"),
page: int = Query(1, ge=1, description="页码,从 1 开始"),
page_size: int = Query(20, ge=1, le=100, description="每页数量,最大 100"),
keyword: str | None = Query(None, description="项目名称模糊搜索"),
status: str | None = Query(None, description="项目状态筛选,不传默认 active"),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
@@ -130,12 +164,22 @@ async def list_private_portrait_projects(
return PrivatePortraitProjectListOut(items=[project_to_out(item) for item in items], total=total, page=page, page_size=page_size)
@router.get("/private-portrait/projects/{project_id}", response_model=PrivatePortraitProjectOut, summary="获取真人认证素材项目详情")
@router.get(
"/private-portrait/projects/{project_id}",
response_model=PrivatePortraitProjectOut,
summary="获取真人认证素材项目详情",
description="获取当前用户真人项目详情,project_id 必须属于当前用户且 library_type=real_person。",
)
async def get_private_portrait_project(project_id: str, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
return project_to_out(await get_user_project(db, user_id=current_user.id, project_id=project_id, library_type=PrivatePortraitLibraryType.REAL_PERSON.value))
@router.put("/private-portrait/projects/{project_id}", response_model=PrivatePortraitProjectOut, summary="更新真人认证素材项目")
@router.put(
"/private-portrait/projects/{project_id}",
response_model=PrivatePortraitProjectOut,
summary="更新真人认证素材项目",
description="更新真人项目本地展示信息;不会重新发起真人认证。每个真人项目只允许认证一次。",
)
async def update_private_portrait_project(project_id: str, payload: PrivatePortraitProjectUpdate, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
project = await update_real_person_project(db, user_id=current_user.id, project_id=project_id, payload=payload)
out = project_to_out(project)
@@ -143,7 +187,12 @@ async def update_private_portrait_project(project_id: str, payload: PrivatePortr
return out
@router.delete("/private-portrait/projects/{project_id}", response_model=PrivatePortraitDeleteOut, summary="删除真人认证素材项目")
@router.delete(
"/private-portrait/projects/{project_id}",
response_model=PrivatePortraitDeleteOut,
summary="删除真人认证素材项目",
description="软删真人项目和本地素材记录。本地先 commit,commit 成功后投递 Celery 删除远端 AssetGroup/Assetremote_delete_status=pending 表示远端删除处理中。",
)
async def delete_private_portrait_project(project_id: str, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
project = await soft_delete_project(db, user_id=current_user.id, project_id=project_id, library_type=PrivatePortraitLibraryType.REAL_PERSON.value)
project_id_snapshot = project.id
@@ -158,7 +207,12 @@ async def delete_private_portrait_project(project_id: str, current_user: User =
return PrivatePortraitDeleteOut(success=True, remote_delete_status=PrivatePortraitRemoteDeleteStatus.PENDING.value)
@router.post("/private-portrait/projects/{project_id}/validate-sessions", response_model=PrivatePortraitValidateSessionOut, summary="重新创建真人认证会话")
@router.post(
"/private-portrait/projects/{project_id}/validate-sessions",
response_model=PrivatePortraitValidateSessionOut,
summary="重新创建真人认证会话",
description="为真人项目创建新的认证会话。业务层会限制项目只能认证一次;已认证成功的项目不允许重复认证。",
)
async def create_private_portrait_validate_session(project_id: str, payload: PrivatePortraitValidateSessionCreate, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
session = await create_real_person_validate_session(db, user_id=current_user.id, project_id=project_id, callback_redirect_url=payload.callback_redirect_url)
out = validate_session_to_out(session)
@@ -166,12 +220,21 @@ async def create_private_portrait_validate_session(project_id: str, payload: Pri
return out
@router.get("/private-portrait/validate-sessions/{session_id}", response_model=PrivatePortraitValidateSessionOut, summary="查询真人认证会话状态")
@router.get(
"/private-portrait/validate-sessions/{session_id}",
response_model=PrivatePortraitValidateSessionOut,
summary="查询真人认证会话状态",
description="PC 端轮询该接口查看手机扫码认证结果。status/result_code/remote_group_id 可用于判断是否认证成功并提示用户回到 PC 查看项目。",
)
async def get_private_portrait_validate_session(session_id: str, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
return validate_session_to_out(await get_validate_session(db, user_id=current_user.id, session_id=session_id))
@router.get("/private-portrait/validate-callback", summary="火山真人认证回调入口")
@router.get(
"/private-portrait/validate-callback",
summary="火山真人认证回调入口",
description="火山真人认证 H5 回调入口。resultCode=10000 表示认证成功;成功后会创建或更新本地 AssetGroup 映射,并可 redirect 回前端提示页。",
)
async def private_portrait_validate_callback(session_id: str, request: Request, redirect_url: str | None = None, db: AsyncSession = Depends(get_db)):
params = dict(request.query_params)
params.pop("session_id", None)
@@ -189,7 +252,16 @@ async def private_portrait_validate_callback(session_id: str, request: Request,
return response
@router.post("/private-portrait/projects/{project_id}/assets", response_model=PrivatePortraitAssetOut, summary="上传真人认证素材", description="当前支持 Image / Video。Audio 暂不开放。CreateAsset 是异步接口,返回后需要轮询到 Active 才可用于生成。")
@router.post(
"/private-portrait/projects/{project_id}/assets",
response_model=PrivatePortraitAssetOut,
summary="上传真人认证素材",
description=(
"在已认证成功的真人项目下创建素材。当前仅开放 asset_type=Image/VideoAudio 暂不开放。"
"Video 必须携带 video_duration,建议前端限制 2~15 秒。CreateAsset 是异步接口,返回后会投递轮询任务,"
"只有 status=Active 的素材才会出现在 selectable-assets 并可用于 AI 创作。"
),
)
async def create_private_portrait_asset(project_id: str, payload: PrivatePortraitAssetCreate, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
asset = await create_real_person_asset(db, user_id=current_user.id, project_id=project_id, payload=payload)
asset_id_snapshot = asset.id
@@ -206,13 +278,32 @@ async def create_private_portrait_asset(project_id: str, payload: PrivatePortrai
return out
@router.get("/private-portrait/projects/{project_id}/assets", response_model=PrivatePortraitAssetListOut, summary="查询真人认证素材列表")
async def list_private_portrait_assets(project_id: str, page: int = Query(1, ge=1), page_size: int = Query(20, ge=1, le=100), status: str | None = Query(None), keyword: str | None = Query(None), asset_type: str | None = Query(None), current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
@router.get(
"/private-portrait/projects/{project_id}/assets",
response_model=PrivatePortraitAssetListOut,
summary="查询真人认证素材列表",
description="查询指定真人项目下的素材。asset_type 可传 Image 或 Videostatus 可筛选素材状态。",
)
async def list_private_portrait_assets(
project_id: str,
page: int = Query(1, ge=1, description="页码,从 1 开始"),
page_size: int = Query(20, ge=1, le=100, description="每页数量,最大 100"),
status: str | None = Query(None, description="素材状态筛选,不传查全部"),
keyword: str | None = Query(None, description="素材名称模糊搜索"),
asset_type: str | None = Query(None, description="素材类型:Image=图片,Video=视频;Audio 暂不开放"),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
assets, total, project_name_map = await list_assets(db, user_id=current_user.id, project_id=project_id, status=status, keyword=keyword, page=page, page_size=page_size, library_type=PrivatePortraitLibraryType.REAL_PERSON.value, asset_type=asset_type)
return PrivatePortraitAssetListOut(items=[asset_to_out(asset, project_name=project_name_map.get(asset.project_id)) for asset in assets], total=total, page=page, page_size=page_size)
@router.get("/private-portrait/assets/{asset_id}", summary="获取真人认证素材详情")
@router.get(
"/private-portrait/assets/{asset_id}",
response_model=PrivatePortraitAssetOut,
summary="获取真人认证素材详情",
description="获取当前用户真人素材详情,asset_id 必须属于当前用户且 library_type=real_person。",
)
async def get_private_portrait_asset(asset_id: str, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
asset = (await db.execute(select(PrivatePortraitAsset).where(PrivatePortraitAsset.id == asset_id, PrivatePortraitAsset.user_id == current_user.id, PrivatePortraitAsset.library_type == PrivatePortraitLibraryType.REAL_PERSON.value).limit(1))).scalar_one_or_none()
if not asset:
@@ -221,7 +312,12 @@ async def get_private_portrait_asset(asset_id: str, current_user: User = Depends
return asset_to_out(asset, project_name=project.name if project else None)
@router.post("/private-portrait/assets/{asset_id}/sync", summary="同步真人认证素材状态")
@router.post(
"/private-portrait/assets/{asset_id}/sync",
response_model=PrivatePortraitAssetOut,
summary="同步真人认证素材状态",
description="主动向火山查询并同步真人素材状态。一般由轮询任务自动执行;前端排查或手动刷新时可调用。",
)
async def sync_private_portrait_asset(asset_id: str, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
asset = await sync_asset_status(db, user_id=current_user.id, asset_id=asset_id)
if asset.library_type != PrivatePortraitLibraryType.REAL_PERSON.value:
@@ -231,7 +327,12 @@ async def sync_private_portrait_asset(asset_id: str, current_user: User = Depend
return out
@router.delete("/private-portrait/assets/{asset_id}", response_model=PrivatePortraitDeleteOut, summary="删除真人认证素材")
@router.delete(
"/private-portrait/assets/{asset_id}",
response_model=PrivatePortraitDeleteOut,
summary="删除真人认证素材",
description="软删本地真人素材记录。本地先 commit,commit 成功后投递 Celery 删除远端 Assetremote_delete_status=pending 表示远端删除处理中。",
)
async def delete_private_portrait_asset(asset_id: str, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
asset = await soft_delete_asset(db, user_id=current_user.id, asset_id=asset_id, library_type=PrivatePortraitLibraryType.REAL_PERSON.value)
asset_id_snapshot = asset.id
@@ -247,7 +348,20 @@ async def delete_private_portrait_asset(asset_id: str, current_user: User = Depe
return PrivatePortraitDeleteOut(success=True, remote_delete_status=PrivatePortraitRemoteDeleteStatus.PENDING.value)
@router.get("/private-portrait/selectable-assets", response_model=PrivatePortraitSelectableAssetListOut, summary="查询可用于生成的真人认证素材")
async def list_private_portrait_selectable_assets(page: int = Query(1, ge=1), page_size: int = Query(20, ge=1, le=100), project_id: str | None = Query(None), keyword: str | None = Query(None), asset_type: str | None = Query(None, description="Image 或 Video,不传查全部。"), current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
@router.get(
"/private-portrait/selectable-assets",
response_model=PrivatePortraitSelectableAssetListOut,
summary="查询可用于 AI 创作的真人认证素材",
description="只返回当前用户真人素材库中 status=Active 且未删除的 Image/Video 素材。该接口给 AI 创作参考内容选择器使用。",
)
async def list_private_portrait_selectable_assets(
page: int = Query(1, ge=1, description="页码,从 1 开始"),
page_size: int = Query(20, ge=1, le=100, description="每页数量,最大 100"),
project_id: str | None = Query(None, description="按真人项目 ID 筛选"),
keyword: str | None = Query(None, description="素材名称模糊搜索"),
asset_type: str | None = Query(None, description="素材类型:Image 或 Video,不传查全部"),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
items, total = await list_selectable_assets(db, user_id=current_user.id, project_id=project_id, keyword=keyword, page=page, page_size=page_size, library_type=PrivatePortraitLibraryType.REAL_PERSON.value, asset_type=asset_type)
return PrivatePortraitSelectableAssetListOut(items=items, total=total, page=page, page_size=page_size)
@@ -50,26 +50,65 @@ from app.services.private_portrait.virtual.service import create_virtual_asset,
router = APIRouter(tags=["私域虚拟人像素材库"])
_VIRTUAL_API_DESCRIPTION = (
"私域虚拟素材库 API。真人和虚拟素材共用用户素材额度;user.private_portrait_asset_limit=0 表示关闭,"
">0 表示启用并限制总素材数量。创建虚拟项目会同步调用火山 CreateAssetGroupGroupType=AIGC"
"remote_project_name 默认使用 default。当前上传只开放 Image/VideoAudio 暂不开放。"
)
def _log_task_dispatch_failed(*, task_name: str, user_id: str | None = None, project_id: str | None = None, asset_id: str | None = None, exc: BaseException) -> None:
log_operation_error(domain=DOMAIN, event_type=PrivatePortraitEventType.TASK_DISPATCH_FAILED.value, source=PrivatePortraitEventSource.API.value, user_id=user_id, project_id=project_id, asset_id=asset_id, exc=exc, detail={"task_name": task_name})
log_operation_error(
domain=DOMAIN,
event_type=PrivatePortraitEventType.TASK_DISPATCH_FAILED.value,
source=PrivatePortraitEventSource.API.value,
user_id=user_id,
project_id=project_id,
asset_id=asset_id,
exc=exc,
detail={"task_name": task_name},
)
def _log_task_dispatch_success(*, task_name: str, user_id: str | None = None, project_id: str | None = None, asset_id: str | None = None) -> None:
log_operation_event(domain=DOMAIN, event_type=PrivatePortraitEventType.TASK_DISPATCH_SUCCESS.value, event_status=PrivatePortraitEventStatus.SUCCESS.value, source=PrivatePortraitEventSource.API.value, user_id=user_id, project_id=project_id, asset_id=asset_id, detail={"task_name": task_name})
log_operation_event(
domain=DOMAIN,
event_type=PrivatePortraitEventType.TASK_DISPATCH_SUCCESS.value,
event_status=PrivatePortraitEventStatus.SUCCESS.value,
source=PrivatePortraitEventSource.API.value,
user_id=user_id,
project_id=project_id,
asset_id=asset_id,
detail={"task_name": task_name},
)
@router.get("/private-portrait/virtual/config", response_model=PrivatePortraitConfigOut, summary="获取虚拟人像素材库额度配置", description="额度与真人认证素材库共用;图片/视频共用;Audio 暂不开放。")
@router.get(
"/private-portrait/virtual/config",
response_model=PrivatePortraitConfigOut,
summary="获取虚拟素材库额度配置",
description=_VIRTUAL_API_DESCRIPTION + "返回虚拟素材库可用额度,实际与真人素材库共用。",
)
async def get_my_virtual_private_portrait_config(current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
return await get_user_private_portrait_config(db, user_id=current_user.id)
@router.get("/private-portrait/virtual-meta/enums", response_model=PrivatePortraitEnumMetaOut, summary="获取虚拟人像素材库枚举说明")
@router.get(
"/private-portrait/virtual-meta/enums",
response_model=PrivatePortraitEnumMetaOut,
summary="获取虚拟素材库枚举说明",
description="返回前端展示所需枚举:library_type、asset_type、project status、asset status、remote_delete_status 等。Audio 暂不开放上传。",
)
async def get_virtual_private_portrait_enum_meta():
return build_private_portrait_enum_meta()
@router.post("/private-portrait/virtual-projects", response_model=PrivatePortraitProjectOut, summary="创建虚拟人像项目组", description="创建本地虚拟人像项目,并同步调用火山 CreateAssetGroupGroupType=AIGC。ProjectName 必须与后续生成 API Key 所属项目一致,默认使用 default。")
@router.post(
"/private-portrait/virtual-projects",
response_model=PrivatePortraitProjectOut,
summary="创建虚拟素材项目组",
description=_VIRTUAL_API_DESCRIPTION + "创建成功后返回本地项目详情,后续上传虚拟素材必须归属到该 project_id。",
)
async def create_private_portrait_virtual_project(payload: PrivatePortraitVirtualProjectCreate, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
project = await create_virtual_project(db, user_id=current_user.id, payload=payload)
out = project_to_out(project)
@@ -77,8 +116,20 @@ async def create_private_portrait_virtual_project(payload: PrivatePortraitVirtua
return out
@router.get("/private-portrait/virtual-projects", response_model=PrivatePortraitProjectListOut, summary="查询当前用户虚拟人像项目列表")
async def list_private_portrait_virtual_projects(page: int = Query(1, ge=1), page_size: int = Query(20, ge=1, le=100), keyword: str | None = Query(None), status: str | None = Query(None), current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
@router.get(
"/private-portrait/virtual-projects",
response_model=PrivatePortraitProjectListOut,
summary="查询当前用户虚拟素材项目列表",
description="只返回 library_type=aigc_virtual 的项目。默认查询 active 项目,可通过 status 覆盖。",
)
async def list_private_portrait_virtual_projects(
page: int = Query(1, ge=1, description="页码,从 1 开始"),
page_size: int = Query(20, ge=1, le=100, description="每页数量,最大 100"),
keyword: str | None = Query(None, description="项目名称模糊搜索"),
status: str | None = Query(None, description="项目状态筛选,不传默认 active"),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
query_status = status or PrivatePortraitProjectStatus.ACTIVE.value
items, total = await list_projects(db, user_id=current_user.id, page=page, page_size=page_size, keyword=keyword, status=query_status, library_type=PrivatePortraitLibraryType.AIGC_VIRTUAL.value)
await refresh_project_counters(db, [item.id for item in items])
@@ -87,12 +138,22 @@ async def list_private_portrait_virtual_projects(page: int = Query(1, ge=1), pag
return PrivatePortraitProjectListOut(items=[project_to_out(item) for item in items], total=total, page=page, page_size=page_size)
@router.get("/private-portrait/virtual-projects/{project_id}", response_model=PrivatePortraitProjectOut, summary="获取虚拟人像项目详情")
@router.get(
"/private-portrait/virtual-projects/{project_id}",
response_model=PrivatePortraitProjectOut,
summary="获取虚拟素材项目详情",
description="获取当前用户虚拟素材项目详情,project_id 必须属于当前用户且 library_type=aigc_virtual。",
)
async def get_private_portrait_virtual_project(project_id: str, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
return project_to_out(await get_user_project(db, user_id=current_user.id, project_id=project_id, library_type=PrivatePortraitLibraryType.AIGC_VIRTUAL.value))
@router.put("/private-portrait/virtual-projects/{project_id}", response_model=PrivatePortraitProjectOut, summary="更新虚拟人像项目")
@router.put(
"/private-portrait/virtual-projects/{project_id}",
response_model=PrivatePortraitProjectOut,
summary="更新虚拟素材项目",
description="更新虚拟素材项目本地展示信息。不会重新创建远端 AssetGroup。",
)
async def update_private_portrait_virtual_project(project_id: str, payload: PrivatePortraitProjectUpdate, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
project = await update_virtual_project(db, user_id=current_user.id, project_id=project_id, payload=payload)
out = project_to_out(project)
@@ -100,7 +161,12 @@ async def update_private_portrait_virtual_project(project_id: str, payload: Priv
return out
@router.delete("/private-portrait/virtual-projects/{project_id}", response_model=PrivatePortraitDeleteOut, summary="删除虚拟人像项目")
@router.delete(
"/private-portrait/virtual-projects/{project_id}",
response_model=PrivatePortraitDeleteOut,
summary="删除虚拟素材项目",
description="软删虚拟素材项目和本地素材记录。本地先 commit,commit 成功后投递 Celery 删除远端 AssetGroup/Assetremote_delete_status=pending 表示远端删除处理中。",
)
async def delete_private_portrait_virtual_project(project_id: str, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
project = await soft_delete_project(db, user_id=current_user.id, project_id=project_id, library_type=PrivatePortraitLibraryType.AIGC_VIRTUAL.value)
project_id_snapshot = project.id
@@ -115,7 +181,16 @@ async def delete_private_portrait_virtual_project(project_id: str, current_user:
return PrivatePortraitDeleteOut(success=True, remote_delete_status=PrivatePortraitRemoteDeleteStatus.PENDING.value)
@router.post("/private-portrait/virtual-projects/{project_id}/assets", response_model=PrivatePortraitAssetOut, summary="上传虚拟人像素材", description="当前支持 Image / Video。Audio 暂不开放。CreateAsset 是异步接口,返回后需要轮询到 Active 才可用于生成。")
@router.post(
"/private-portrait/virtual-projects/{project_id}/assets",
response_model=PrivatePortraitAssetOut,
summary="上传虚拟素材",
description=(
"在虚拟素材项目下创建素材。当前仅开放 asset_type=Image/VideoAudio 暂不开放。"
"Video 必须携带 video_duration,建议前端限制 2~15 秒。CreateAsset 是异步接口,返回后会投递轮询任务,"
"只有 status=Active 的素材才会出现在 virtual-selectable-assets 并可用于 AI 创作。"
),
)
async def create_private_portrait_virtual_asset(project_id: str, payload: PrivatePortraitAssetCreate, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
asset = await create_virtual_asset(db, user_id=current_user.id, project_id=project_id, payload=payload)
asset_id_snapshot = asset.id
@@ -132,13 +207,32 @@ async def create_private_portrait_virtual_asset(project_id: str, payload: Privat
return out
@router.get("/private-portrait/virtual-projects/{project_id}/assets", response_model=PrivatePortraitAssetListOut, summary="查询虚拟人像素材列表")
async def list_private_portrait_virtual_assets(project_id: str, page: int = Query(1, ge=1), page_size: int = Query(20, ge=1, le=100), status: str | None = Query(None), keyword: str | None = Query(None), asset_type: str | None = Query(None), current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
@router.get(
"/private-portrait/virtual-projects/{project_id}/assets",
response_model=PrivatePortraitAssetListOut,
summary="查询虚拟素材列表",
description="查询指定虚拟项目下的素材。asset_type 可传 Image 或 Videostatus 可筛选素材状态。",
)
async def list_private_portrait_virtual_assets(
project_id: str,
page: int = Query(1, ge=1, description="页码,从 1 开始"),
page_size: int = Query(20, ge=1, le=100, description="每页数量,最大 100"),
status: str | None = Query(None, description="素材状态筛选,不传查全部"),
keyword: str | None = Query(None, description="素材名称模糊搜索"),
asset_type: str | None = Query(None, description="素材类型:Image=图片,Video=视频;Audio 暂不开放"),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
assets, total, project_name_map = await list_assets(db, user_id=current_user.id, project_id=project_id, status=status, keyword=keyword, page=page, page_size=page_size, library_type=PrivatePortraitLibraryType.AIGC_VIRTUAL.value, asset_type=asset_type)
return PrivatePortraitAssetListOut(items=[asset_to_out(asset, project_name=project_name_map.get(asset.project_id)) for asset in assets], total=total, page=page, page_size=page_size)
@router.get("/private-portrait/virtual-assets/{asset_id}", response_model=PrivatePortraitAssetOut, summary="获取虚拟人像素材详情")
@router.get(
"/private-portrait/virtual-assets/{asset_id}",
response_model=PrivatePortraitAssetOut,
summary="获取虚拟素材详情",
description="获取当前用户虚拟素材详情,asset_id 必须属于当前用户且 library_type=aigc_virtual。",
)
async def get_private_portrait_virtual_asset(asset_id: str, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
asset = (await db.execute(select(PrivatePortraitAsset).where(PrivatePortraitAsset.id == asset_id, PrivatePortraitAsset.user_id == current_user.id, PrivatePortraitAsset.library_type == PrivatePortraitLibraryType.AIGC_VIRTUAL.value).limit(1))).scalar_one_or_none()
if not asset:
@@ -147,7 +241,12 @@ async def get_private_portrait_virtual_asset(asset_id: str, current_user: User =
return asset_to_out(asset, project_name=project.name if project else None)
@router.post("/private-portrait/virtual-assets/{asset_id}/sync", response_model=PrivatePortraitAssetOut, summary="同步虚拟人像素材状态")
@router.post(
"/private-portrait/virtual-assets/{asset_id}/sync",
response_model=PrivatePortraitAssetOut,
summary="同步虚拟素材状态",
description="主动向火山查询并同步虚拟素材状态。一般由轮询任务自动执行;前端排查或手动刷新时可调用。",
)
async def sync_private_portrait_virtual_asset(asset_id: str, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
asset = await sync_asset_status(db, user_id=current_user.id, asset_id=asset_id)
if asset.library_type != PrivatePortraitLibraryType.AIGC_VIRTUAL.value:
@@ -157,7 +256,12 @@ async def sync_private_portrait_virtual_asset(asset_id: str, current_user: User
return out
@router.delete("/private-portrait/virtual-assets/{asset_id}", response_model=PrivatePortraitDeleteOut, summary="删除虚拟人像素材")
@router.delete(
"/private-portrait/virtual-assets/{asset_id}",
response_model=PrivatePortraitDeleteOut,
summary="删除虚拟素材",
description="软删本地虚拟素材记录。本地先 commit,commit 成功后投递 Celery 删除远端 Assetremote_delete_status=pending 表示远端删除处理中。",
)
async def delete_private_portrait_virtual_asset(asset_id: str, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
asset = await soft_delete_asset(db, user_id=current_user.id, asset_id=asset_id, library_type=PrivatePortraitLibraryType.AIGC_VIRTUAL.value)
asset_id_snapshot = asset.id
@@ -173,7 +277,20 @@ async def delete_private_portrait_virtual_asset(asset_id: str, current_user: Use
return PrivatePortraitDeleteOut(success=True, remote_delete_status=PrivatePortraitRemoteDeleteStatus.PENDING.value)
@router.get("/private-portrait/virtual-selectable-assets", response_model=PrivatePortraitSelectableAssetListOut, summary="查询可用于生成的虚拟人像素材")
async def list_private_portrait_virtual_selectable_assets(page: int = Query(1, ge=1), page_size: int = Query(20, ge=1, le=100), project_id: str | None = Query(None), keyword: str | None = Query(None), asset_type: str | None = Query(None, description="Image 或 Video,不传查全部。"), current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)):
@router.get(
"/private-portrait/virtual-selectable-assets",
response_model=PrivatePortraitSelectableAssetListOut,
summary="查询可用于 AI 创作的虚拟素材",
description="只返回当前用户虚拟素材库中 status=Active 且未删除的 Image/Video 素材。该接口给 AI 创作参考内容选择器使用。",
)
async def list_private_portrait_virtual_selectable_assets(
page: int = Query(1, ge=1, description="页码,从 1 开始"),
page_size: int = Query(20, ge=1, le=100, description="每页数量,最大 100"),
project_id: str | None = Query(None, description="按虚拟项目 ID 筛选"),
keyword: str | None = Query(None, description="素材名称模糊搜索"),
asset_type: str | None = Query(None, description="素材类型:Image 或 Video,不传查全部"),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
items, total = await list_selectable_assets(db, user_id=current_user.id, project_id=project_id, keyword=keyword, page=page, page_size=page_size, library_type=PrivatePortraitLibraryType.AIGC_VIRTUAL.value, asset_type=asset_type)
return PrivatePortraitSelectableAssetListOut(items=items, total=total, page=page, page_size=page_size)
+197 -21
View File
@@ -3,7 +3,7 @@ from __future__ import annotations
from datetime import datetime
from types import SimpleNamespace
from fastapi import APIRouter, Body, Depends, HTTPException, Path, Query
from fastapi import APIRouter, Body, Depends, File, HTTPException, Path, Query, UploadFile
from sqlalchemy import inspect as sa_inspect
from sqlalchemy.ext.asyncio import AsyncSession
@@ -35,6 +35,7 @@ from app.schemas.shot_replicate import (
ShotReplicateTaskDetailOut,
ShotReplicateVideoPromptSchemaUpdateRequest,
ShotSegmentDeleteOut,
ShotTaskSetDeleteOut,
ShotSegmentDetailOut,
ShotSegmentListOut,
ShotSegmentReplicationCreateRequest,
@@ -49,7 +50,6 @@ from app.schemas.shot_replicate import (
from app.services.shot_replicate_flow_service import (
_get_project_for_user,
create_shot_replicate_project_from_segment,
delete_shot_replicate_project,
generate_image_from_prompt,
generate_video_from_prompt,
mark_shot_replicate_step_dispatch_failed,
@@ -65,6 +65,7 @@ from app.services.shot_replicate_taskset_service import (
create_segments_by_ai,
create_task_set,
delete_segment,
delete_task_set,
get_segment_for_user,
list_segments,
list_task_sets,
@@ -83,6 +84,9 @@ from app.services.module_async_recovery_service import (
register_shot_task_set_analysis_task,
)
from app.tasks.celery_app import celery_app
from app.enums.upload_resource import UploadResourceEventEnum, UploadResourceModuleEnum, UploadResourceSourceModelEnum, UploadResourceTypeEnum
from app.services.upload_resource import upload_reference_file, bind_upload_resources, cleanup_upload_resource_files_after_commit
from app.services.upload_resource.log_service import log_upload_resource_exception, safe_rollback_with_log
MODULE = ModuleCodeEnum.SHOT_REPLICATE.value
@@ -245,6 +249,66 @@ async def get_spec():
return ShotReplicateSpecOut()
@router.post(
"/upload-image",
summary="上传拆镜复刻图片素材",
description="上传后先记录为 pending UploadResource;创建拆镜复刻项目后绑定业务记录。",
)
async def upload_shot_replicate_image(
file: UploadFile = File(...),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
result = await upload_reference_file(
db,
file=file,
current_user=current_user,
module=UploadResourceModuleEnum.SHOT_REPLICATE.value,
resource_type=UploadResourceTypeEnum.IMAGE.value,
gen_type="video",
)
await db.commit()
return {
"url": result.url,
"filename": result.filename,
"type": "image",
"module": result.module,
"resource_id": result.resource_id,
"file_size_bytes": result.file_size_bytes,
}
@router.post(
"/upload-video",
summary="上传拆镜复刻视频素材",
description="上传后先记录为 pending UploadResource;创建拆镜总任务集后绑定 ShotReplicateTaskSet。",
)
async def upload_shot_replicate_video(
file: UploadFile = File(...),
duration_seconds: float | None = Query(None, description="前端识别的视频时长秒数,可选"),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
result = await upload_reference_file(
db,
file=file,
current_user=current_user,
module=UploadResourceModuleEnum.SHOT_REPLICATE.value,
resource_type=UploadResourceTypeEnum.VIDEO.value,
duration_seconds=duration_seconds,
)
await db.commit()
return {
"url": result.url,
"filename": result.filename,
"type": "video",
"module": result.module,
"resource_id": result.resource_id,
"file_size_bytes": result.file_size_bytes,
"duration_seconds": result.duration_seconds,
}
@router.post(
"/task-sets",
response_model=ShotTaskSetDetailOut,
@@ -264,6 +328,16 @@ async def create_shot_task_set(
try:
task_set = await create_task_set(db, current_user=current_user, req=req)
task_set_id = task_set.id
await bind_upload_resources(
db,
user_id=current_user.id,
module=UploadResourceModuleEnum.SHOT_REPLICATE.value,
source_model=UploadResourceSourceModelEnum.SHOT_REPLICATE_TASK_SET.value,
source_id=task_set_id,
resource_ids=[req.video_resource_id],
urls=[req.video_url],
allow_common_migrate=True,
)
await db.commit()
except HTTPException:
await db.rollback()
@@ -626,18 +700,69 @@ async def delete_shot_segment(
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
user_id = _safe_user_id(current_user)
pending_ids: list[str] = []
try:
out = await delete_segment(db, current_user=current_user, segment_id=segment_id)
pending_ids = list(out.pending_delete_resource_ids or [])
await db.commit()
return out
except HTTPException:
await db.rollback()
except HTTPException as exc:
await safe_rollback_with_log(
db,
event_type=UploadResourceEventEnum.SHOT_SEGMENT_DELETE_ROLLBACK_FAILED.value,
message="删除拆镜片段 HTTPException 回滚失败",
user_id=user_id,
module=MODULE,
detail={"segment_id": segment_id},
original_exc=exc,
)
raise
except Exception as exc:
await db.rollback()
except Exception as exc: # noqa: BLE001
await safe_rollback_with_log(
db,
event_type=UploadResourceEventEnum.SHOT_SEGMENT_DELETE_ROLLBACK_FAILED.value,
message="删除拆镜片段主事务回滚失败",
user_id=user_id,
module=MODULE,
detail={"segment_id": segment_id},
original_exc=exc,
)
_log_api_exception_from_locals(exc, locals(), f"删除拆镜片段失败: {exc}")
log_upload_resource_exception(
event_type=UploadResourceEventEnum.SHOT_SEGMENT_DELETE_FAILED.value,
message=f"删除拆镜片段失败: {exc}",
user_id=user_id,
module=MODULE,
detail={"segment_id": segment_id},
exc=exc,
)
raise HTTPException(status_code=500, detail=f"删除拆镜片段失败: {exc}")
if pending_ids:
try:
await cleanup_upload_resource_files_after_commit(db, resource_ids=pending_ids)
await db.commit()
except Exception as exc: # noqa: BLE001
await safe_rollback_with_log(
db,
event_type=UploadResourceEventEnum.SHOT_SEGMENT_DELETE_ROLLBACK_FAILED.value,
message="删除拆镜片段 cleanup 事务回滚失败",
user_id=user_id,
module=MODULE,
detail={"segment_id": segment_id, "pending_ids": pending_ids},
original_exc=exc,
)
log_upload_resource_exception(
event_type=UploadResourceEventEnum.SHOT_SEGMENT_DELETE_CLEANUP_FAILED.value,
message=f"删除拆镜片段后清理真实文件失败: {exc}",
user_id=user_id,
resource_ids=pending_ids,
module=MODULE,
detail={"segment_id": segment_id},
exc=exc,
)
return out
@router.post(
"/segments/{segment_id}/replication-projects",
@@ -950,24 +1075,75 @@ async def generate_video(
@router.delete(
"/projects/{project_id}",
response_model=ShotReplicateDeleteOut,
summary="软删除拆镜复刻项目",
description="软删除拆镜复刻 ModuleGenerationProject,并联动软删除当前有效步骤和关联的 ChatGenerationTask",
"/task-sets/{task_set_id}",
response_model=ShotTaskSetDeleteOut,
summary="软删除拆镜任务集",
description="软删除整个 ShotReplicateTaskSet,并联动软删除全部片段和片段内部复刻项目;UploadResource 真实文件在事务提交后清理",
)
async def delete_project(
project_id: str = Path(..., description="拆镜复刻项目ID,即 module_generation_projects.id"),
async def delete_shot_task_set(
task_set_id: str = Path(..., description="拆镜任务集ID,即 shot_replicate_task_sets.id"),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
user_id = _safe_user_id(current_user)
pending_ids: list[str] = []
try:
out = await delete_shot_replicate_project(db, current_user=current_user, project_id=project_id)
out = await delete_task_set(db, current_user=current_user, task_set_id=task_set_id)
pending_ids = list(out.pending_delete_resource_ids or [])
await db.commit()
return out
except HTTPException:
await db.rollback()
except HTTPException as exc:
await safe_rollback_with_log(
db,
event_type=UploadResourceEventEnum.SHOT_TASK_SET_DELETE_ROLLBACK_FAILED.value,
message="删除拆镜任务集 HTTPException 回滚失败",
user_id=user_id,
module=MODULE,
detail={"task_set_id": task_set_id},
original_exc=exc,
)
raise
except Exception as exc:
await db.rollback()
_log_api_exception_from_locals(exc, locals(), f"删除拆镜复刻项目失败: {exc}")
raise HTTPException(status_code=500, detail=f"删除拆镜复刻项目失败: {exc}")
except Exception as exc: # noqa: BLE001
await safe_rollback_with_log(
db,
event_type=UploadResourceEventEnum.SHOT_TASK_SET_DELETE_ROLLBACK_FAILED.value,
message="删除拆镜任务集主事务回滚失败",
user_id=user_id,
module=MODULE,
detail={"task_set_id": task_set_id},
original_exc=exc,
)
_log_api_exception_from_locals(exc, locals(), f"删除拆镜任务集失败: {exc}")
log_upload_resource_exception(
event_type=UploadResourceEventEnum.SHOT_TASK_SET_DELETE_FAILED.value,
message=f"删除拆镜任务集失败: {exc}",
user_id=user_id,
module=MODULE,
detail={"task_set_id": task_set_id},
exc=exc,
)
raise HTTPException(status_code=500, detail=f"删除拆镜任务集失败: {exc}")
if pending_ids:
try:
await cleanup_upload_resource_files_after_commit(db, resource_ids=pending_ids)
await db.commit()
except Exception as exc: # noqa: BLE001
await safe_rollback_with_log(
db,
event_type=UploadResourceEventEnum.SHOT_TASK_SET_DELETE_ROLLBACK_FAILED.value,
message="删除拆镜任务集 cleanup 事务回滚失败",
user_id=user_id,
module=MODULE,
detail={"task_set_id": task_set_id, "pending_ids": pending_ids},
original_exc=exc,
)
log_upload_resource_exception(
event_type=UploadResourceEventEnum.SHOT_TASK_SET_DELETE_CLEANUP_FAILED.value,
message=f"删除拆镜任务集后清理真实文件失败: {exc}",
user_id=user_id,
resource_ids=pending_ids,
module=MODULE,
detail={"task_set_id": task_set_id},
exc=exc,
)
return out
+225
View File
@@ -0,0 +1,225 @@
from __future__ import annotations
from fastapi import APIRouter, Body, Depends, HTTPException, Path, Query
from sqlalchemy.ext.asyncio import AsyncSession
from app.dependencies import get_current_user, get_db
from app.enums.upload_resource import UploadResourceEventEnum
from app.models.user import User
from app.schemas.upload_resource import (
UploadResourceHistoryBatchDeleteOut,
UploadResourceHistoryBatchDeleteRequest,
UploadResourceHistoryDayItemsOut,
UploadResourceHistoryGroupedOut,
)
from app.services.upload_resource.delete_service import (
cleanup_upload_resource_history_files,
mark_upload_resource_history_deleted,
)
from app.services.upload_resource.history_service import (
list_upload_resource_history_day_items,
list_upload_resource_history_grouped_days,
)
from app.services.upload_resource.log_service import (
log_upload_resource_exception,
log_upload_resource_event,
safe_rollback_with_log,
)
router = APIRouter(
prefix="/upload-resources",
tags=["upload-resources"],
)
_HISTORY_DESCRIPTION = (
"查询当前用户可展示、可复用、可删除的上传素材历史。"
"只返回 UploadResource 中未删除、未绑定业务记录、bind_status=pending、"
"delete_policy=user_deletable、resource_type=image/video/audio 的记录。"
"已绑定爆款开头复刻、拆镜复刻、拆镜切片或其它业务记录的素材不会返回,也不能通过该模块删除。"
)
@router.get(
"/history",
response_model=UploadResourceHistoryGroupedOut,
summary="获取上传历史素材日期分组",
description=(
_HISTORY_DESCRIPTION
+ "返回结构与 /generation-ai/history 接近,generated_date 表示上传日期。"
"每个日期默认回填前 10 条 items,避免前端首屏再次逐日请求。"
),
responses={
200: {"description": "查询成功,返回上传素材日期分组"},
400: {"description": "resource_type 不合法或分页参数不合法"},
401: {"description": "未登录或 Token 无效"},
},
)
async def list_history_grouped_days(
resource_type: str | None = Query(
None,
description="素材类型筛选:image=图片,video=视频,audio=音频;为空表示全部",
examples=["image"],
),
page: int = Query(1, ge=1, description="日期分组页码,从 1 开始"),
page_size: int = Query(10, ge=1, le=10, description="日期分组每页数量,最大 10"),
keyword: str | None = Query(None, description="文件名或资源 URL 模糊搜索,可为空"),
scene: str | None = Query(
None,
description="前端使用场景标识:record=素材云展示,picker=AI创作选择器;仅用于日志排查,不影响过滤规则",
examples=["picker"],
),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
try:
return await list_upload_resource_history_grouped_days(
db,
user_id=current_user.id,
resource_type=resource_type,
page=page,
page_size=page_size,
keyword=keyword,
)
except Exception as exc: # noqa: BLE001
log_upload_resource_exception(
event_type=UploadResourceEventEnum.HISTORY_LIST_FAILED.value,
message=f"查询上传历史素材日期分组失败: {exc}",
user_id=current_user.id,
detail={"resource_type": resource_type, "page": page, "page_size": page_size, "keyword": keyword, "scene": scene},
exc=exc,
)
raise
@router.get(
"/history/{generated_date}",
response_model=UploadResourceHistoryDayItemsOut,
summary="获取指定日期下上传历史素材",
description=(
_HISTORY_DESCRIPTION
+ "generated_date 格式为 YYYY-MM-DD,字段名沿用生成历史接口,实际语义为上传日期。"
),
responses={
200: {"description": "查询成功,返回指定上传日期下的素材列表"},
400: {"description": "日期格式错误、resource_type 不合法或分页参数不合法"},
401: {"description": "未登录或 Token 无效"},
},
)
async def list_history_day_items(
generated_date: str = Path(..., description="上传日期,格式 YYYY-MM-DD", examples=["2026-07-08"]),
resource_type: str | None = Query(
None,
description="素材类型筛选:image=图片,video=视频,audio=音频;为空表示全部",
examples=["video"],
),
page: int = Query(1, ge=1, description="当前日期下分页页码,从 1 开始"),
page_size: int = Query(20, ge=1, le=100, description="当前日期下每页数量,最大 100"),
keyword: str | None = Query(None, description="文件名或资源 URL 模糊搜索,可为空"),
scene: str | None = Query(
None,
description="前端使用场景标识:record=素材云展示,picker=AI创作选择器;仅用于日志排查,不影响过滤规则",
examples=["record"],
),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
try:
return await list_upload_resource_history_day_items(
db,
user_id=current_user.id,
generated_date=generated_date,
resource_type=resource_type,
page=page,
page_size=page_size,
keyword=keyword,
)
except Exception as exc: # noqa: BLE001
log_upload_resource_exception(
event_type=UploadResourceEventEnum.HISTORY_DAY_LIST_FAILED.value,
message=f"查询指定日期上传历史素材失败: {exc}",
user_id=current_user.id,
detail={"generated_date": generated_date, "resource_type": resource_type, "page": page, "page_size": page_size, "keyword": keyword, "scene": scene},
exc=exc,
)
raise
@router.delete(
"/history/batch",
response_model=UploadResourceHistoryBatchDeleteOut,
summary="批量删除上传历史素材",
description=(
"批量删除当前用户的上传历史素材。一次最多 30 条。"
"只允许删除未绑定业务记录、bind_status=pending、delete_policy=user_deletable 的 UploadResource。"
"删除顺序为:主事务内软删 DB 记录并释放上传容量 -> commit 成功 -> 再删除真实静态文件 -> 写回真实文件删除状态。"
"如果真实文件删除失败,主删除不回滚,file_delete_status 会保留 delete_failed/pending_delete 供后续补偿排查。"
),
responses={
200: {"description": "删除成功,返回软删数量、释放容量和真实文件清理结果"},
400: {"description": "resource_ids 为空、重复或超过 30 条"},
404: {"description": "部分素材不存在、已删除、已绑定业务记录或无权操作"},
401: {"description": "未登录或 Token 无效"},
},
)
async def batch_delete_upload_history(
req: UploadResourceHistoryBatchDeleteRequest = Body(..., description="批量删除上传历史素材请求体"),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
result: UploadResourceHistoryBatchDeleteOut | None = None
user_id = str(current_user.id)
try:
result = await mark_upload_resource_history_deleted(db, current_user=current_user, resource_ids=req.resource_ids)
await db.commit()
log_upload_resource_event(
event_type=UploadResourceEventEnum.DELETE_BATCH_COMMIT_SUCCESS.value,
user_id=user_id,
detail={
"requested_ids": result.requested_ids,
"deleted_ids": result.deleted_ids,
"released_size_bytes": result.released_size_bytes,
},
)
except Exception as exc: # noqa: BLE001
await safe_rollback_with_log(
db,
event_type=UploadResourceEventEnum.DELETE_BATCH_ROLLBACK_FAILED.value,
message="批量删除上传历史素材主事务回滚失败",
user_id=user_id,
detail={"resource_ids": req.resource_ids},
original_exc=exc,
)
log_upload_resource_exception(
event_type=UploadResourceEventEnum.DELETE_BATCH_FAILED.value,
message=f"批量删除上传历史素材失败: {exc}",
user_id=user_id,
resource_ids=req.resource_ids,
detail={"resource_ids": req.resource_ids},
exc=exc,
)
raise
if not result:
raise HTTPException(status_code=500, detail="批量删除上传历史素材失败")
try:
result = await cleanup_upload_resource_history_files(db, result=result, current_user=current_user)
await db.commit()
except Exception as exc: # noqa: BLE001
await safe_rollback_with_log(
db,
event_type=UploadResourceEventEnum.DELETE_BATCH_ROLLBACK_FAILED.value,
message="批量删除上传历史素材 cleanup 事务回滚失败",
user_id=user_id,
detail={"resource_ids": req.resource_ids, "deleted_ids": result.deleted_ids if result else []},
original_exc=exc,
)
log_upload_resource_exception(
event_type=UploadResourceEventEnum.DELETE_BATCH_CLEANUP_FAILED.value,
message=f"批量删除上传历史素材后清理真实文件失败: {exc}",
user_id=user_id,
resource_ids=result.deleted_ids if result else req.resource_ids,
detail={"resource_ids": req.resource_ids, "deleted_ids": result.deleted_ids if result else []},
exc=exc,
)
return result