拆镜复刻片段重试切片API | 真人/虚拟素材库上传API
This commit is contained in:
@@ -6,7 +6,7 @@ from typing import Any
|
||||
from urllib.parse import urlencode
|
||||
|
||||
from fastapi import HTTPException
|
||||
from sqlalchemy import func, select
|
||||
from sqlalchemy import func, or_, select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.config import settings
|
||||
@@ -31,11 +31,22 @@ from app.enums.private_portrait import (
|
||||
PrivatePortraitRemoteDeleteStatus,
|
||||
PrivatePortraitValidateSessionStatus,
|
||||
)
|
||||
from app.enums.upload_resource import (
|
||||
UploadResourceBindStatusEnum,
|
||||
UploadResourceDeletePolicyEnum,
|
||||
UploadResourceModuleEnum,
|
||||
UploadResourceSourceModelEnum,
|
||||
UploadResourceTypeEnum,
|
||||
)
|
||||
from app.models.private_portrait import PrivatePortraitAsset, PrivatePortraitAssetGroup, PrivatePortraitProject, PrivatePortraitValidateSession
|
||||
from app.models.upload_resource import UploadResource
|
||||
from app.schemas.private_portrait import PrivatePortraitAssetCreate, PrivatePortraitAssetOut, PrivatePortraitSelectableAssetOut, PrivatePortraitValidateSessionOut
|
||||
from app.services.operation_log_service import log_operation_error, log_operation_event
|
||||
from app.services.private_portrait.ark_client import ArkPrivateAssetClient
|
||||
from app.services.private_portrait.project_service import get_user_project, refresh_project_counters
|
||||
from app.services.private_portrait.upload_service import private_portrait_upload_module
|
||||
from app.services.upload_resource import bind_upload_resources, release_upload_resources_by_source
|
||||
from app.services.upload_resource.path_resolver import upload_url_to_storage_path
|
||||
from app.services.private_portrait.quota_service import (
|
||||
count_user_counting_assets,
|
||||
ensure_private_portrait_asset_quota_available,
|
||||
@@ -133,6 +144,78 @@ def _assert_private_asset_video_duration(payload: PrivatePortraitAssetCreate) ->
|
||||
raise HTTPException(status_code=400, detail=f"视频素材最长不能超过 {PRIVATE_PORTRAIT_VIDEO_MAX_DURATION_SECONDS} 秒")
|
||||
|
||||
|
||||
def _resource_type_for_asset_type(asset_type: str) -> str:
|
||||
if asset_type == PrivatePortraitAssetType.VIDEO.value:
|
||||
return UploadResourceTypeEnum.VIDEO.value
|
||||
return UploadResourceTypeEnum.IMAGE.value
|
||||
|
||||
|
||||
def _safe_set_payload_attr(payload: PrivatePortraitAssetCreate, name: str, value: Any) -> None:
|
||||
if value is None:
|
||||
return
|
||||
try:
|
||||
setattr(payload, name, value)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
async def _resolve_upload_resource_for_asset(
|
||||
db: AsyncSession,
|
||||
*,
|
||||
user_id: str,
|
||||
payload: PrivatePortraitAssetCreate,
|
||||
module: str,
|
||||
) -> UploadResource | None:
|
||||
"""定位并校验待绑定的 UploadResource。
|
||||
|
||||
新客户端必须传 upload_resource_id;旧客户端没传时按 url 反查 storage_path 兼容。
|
||||
不通过 relationship 懒加载,全部按 ID/路径批查,避免 commit 后 ORM 失效风险。
|
||||
"""
|
||||
resource_id = str(payload.upload_resource_id or "").strip() or None
|
||||
storage_path = upload_url_to_storage_path(payload.url)
|
||||
if not resource_id and not storage_path:
|
||||
return None
|
||||
|
||||
filters = [
|
||||
UploadResource.user_id == user_id,
|
||||
UploadResource.deleted_at.is_(None),
|
||||
]
|
||||
if resource_id and storage_path:
|
||||
filters.append(or_(UploadResource.id == resource_id, UploadResource.storage_path == storage_path))
|
||||
elif resource_id:
|
||||
filters.append(UploadResource.id == resource_id)
|
||||
else:
|
||||
filters.append(UploadResource.storage_path == storage_path)
|
||||
|
||||
result = await db.execute(select(UploadResource).where(*filters).with_for_update().limit(1))
|
||||
resource = result.scalar_one_or_none()
|
||||
if not resource:
|
||||
if resource_id:
|
||||
raise HTTPException(status_code=404, detail="上传资源不存在或不属于当前用户")
|
||||
return None
|
||||
|
||||
if resource.bind_status != UploadResourceBindStatusEnum.PENDING.value or resource.source_id or resource.source_model:
|
||||
raise HTTPException(status_code=409, detail="上传资源已绑定其他素材,不能重复使用")
|
||||
if resource.delete_policy != UploadResourceDeletePolicyEnum.USER_DELETABLE.value:
|
||||
raise HTTPException(status_code=409, detail="上传资源当前不允许绑定私域素材")
|
||||
|
||||
expected_type = _resource_type_for_asset_type(payload.asset_type)
|
||||
if resource.resource_type != expected_type:
|
||||
raise HTTPException(status_code=400, detail="上传资源类型与素材类型不一致")
|
||||
|
||||
if resource.module not in {module, UploadResourceModuleEnum.COMMON.value}:
|
||||
raise HTTPException(status_code=409, detail="上传资源所属模块不匹配,请重新上传素材")
|
||||
|
||||
if payload.asset_type == PrivatePortraitAssetType.VIDEO.value and payload.video_duration is None and resource.duration_seconds is not None:
|
||||
_safe_set_payload_attr(payload, "video_duration", float(resource.duration_seconds))
|
||||
if payload.file_size is None and resource.file_size_bytes is not None:
|
||||
_safe_set_payload_attr(payload, "file_size", int(resource.file_size_bytes or 0))
|
||||
if payload.mime_type is None and resource.mime_type:
|
||||
_safe_set_payload_attr(payload, "mime_type", resource.mime_type)
|
||||
|
||||
return resource
|
||||
|
||||
|
||||
def validate_session_to_out(session: PrivatePortraitValidateSession, *, include_user: bool = False) -> PrivatePortraitValidateSessionOut:
|
||||
return PrivatePortraitValidateSessionOut(
|
||||
id=session.id,
|
||||
@@ -417,11 +500,14 @@ async def create_asset(
|
||||
library_type: str | None = None,
|
||||
) -> PrivatePortraitAsset:
|
||||
_assert_enabled_asset_type(payload.asset_type)
|
||||
_assert_private_asset_video_duration(payload)
|
||||
project = await get_user_project(db, user_id=user_id, project_id=project_id, library_type=library_type)
|
||||
if project.status != PrivatePortraitProjectStatus.ACTIVE.value:
|
||||
raise HTTPException(status_code=400, detail="项目未激活,不能上传素材")
|
||||
|
||||
module = private_portrait_upload_module(project.library_type)
|
||||
upload_resource = await _resolve_upload_resource_for_asset(db, user_id=user_id, payload=payload, module=module)
|
||||
_assert_private_asset_video_duration(payload)
|
||||
|
||||
limit, current_count = await ensure_private_portrait_asset_quota_available(db, user_id=user_id, project_id=project_id, library_type=project.library_type, asset_type=payload.asset_type)
|
||||
group = await get_project_active_group(db, user_id=user_id, project_id=project.id, library_type=project.library_type)
|
||||
public_url = _public_url(payload.url)
|
||||
@@ -445,7 +531,7 @@ async def create_asset(
|
||||
)
|
||||
db.add(asset)
|
||||
await db.flush()
|
||||
log_operation_event(domain=DOMAIN, event_type=PrivatePortraitEventType.ASSET_CREATE_START.value, event_status=PrivatePortraitEventStatus.PENDING.value, source=PrivatePortraitEventSource.API.value, user_id=user_id, project_id=project.id, group_id=group.id, asset_id=asset.id, detail={"asset_limit": limit, "used_asset_count": current_count, "library_type": project.library_type, "asset_type": payload.asset_type, "remote_project_name": project.remote_project_name})
|
||||
log_operation_event(domain=DOMAIN, event_type=PrivatePortraitEventType.ASSET_CREATE_START.value, event_status=PrivatePortraitEventStatus.PENDING.value, source=PrivatePortraitEventSource.API.value, user_id=user_id, project_id=project.id, group_id=group.id, asset_id=asset.id, detail={"asset_limit": limit, "used_asset_count": current_count, "library_type": project.library_type, "asset_type": payload.asset_type, "remote_project_name": project.remote_project_name, "upload_resource_id": upload_resource.id if upload_resource else payload.upload_resource_id})
|
||||
try:
|
||||
remote_resp = await ArkPrivateAssetClient().create_asset(project_name=project.remote_project_name, group_id=group.remote_group_id, url=public_url, asset_type=payload.asset_type, name=payload.name)
|
||||
remote_asset_id = remote_resp.get("Id") or remote_resp.get("AssetId") or remote_resp.get("assetId")
|
||||
@@ -456,6 +542,28 @@ async def create_asset(
|
||||
asset.status = PrivatePortraitAssetStatus.PROCESSING.value
|
||||
asset.next_poll_at = now + timedelta(seconds=_poll_interval_seconds(asset.asset_type))
|
||||
asset.raw_response_json = _json(remote_resp)
|
||||
bind_stats = await bind_upload_resources(
|
||||
db,
|
||||
user_id=user_id,
|
||||
module=module,
|
||||
source_model=UploadResourceSourceModelEnum.PRIVATE_PORTRAIT_ASSET.value,
|
||||
source_id=asset.id,
|
||||
resource_ids=[payload.upload_resource_id, upload_resource.id if upload_resource else None],
|
||||
urls=[payload.url],
|
||||
allow_common_migrate=True,
|
||||
)
|
||||
if bind_stats.get("bound"):
|
||||
log_operation_event(
|
||||
domain=DOMAIN,
|
||||
event_type=PrivatePortraitEventType.ASSET_UPLOAD_BIND_SUCCESS.value,
|
||||
event_status=PrivatePortraitEventStatus.SUCCESS.value,
|
||||
source=PrivatePortraitEventSource.API.value,
|
||||
user_id=user_id,
|
||||
project_id=project.id,
|
||||
group_id=group.id,
|
||||
asset_id=asset.id,
|
||||
detail={"module": module, "upload_resource_id": payload.upload_resource_id, "bind_stats": bind_stats},
|
||||
)
|
||||
await refresh_project_counters(db, [project.id])
|
||||
await db.flush()
|
||||
await db.refresh(asset)
|
||||
@@ -465,7 +573,7 @@ async def create_asset(
|
||||
asset.status = PrivatePortraitAssetStatus.FAILED.value
|
||||
asset.error_message = _exception_message(exc)
|
||||
await db.flush()
|
||||
log_operation_error(domain=DOMAIN, event_type=PrivatePortraitEventType.ASSET_CREATE_FAILED.value, source=PrivatePortraitEventSource.API.value, user_id=user_id, project_id=project.id, group_id=group.id, asset_id=asset.id, exc=exc)
|
||||
log_operation_error(domain=DOMAIN, event_type=PrivatePortraitEventType.ASSET_CREATE_FAILED.value, source=PrivatePortraitEventSource.API.value, user_id=user_id, project_id=project.id, group_id=group.id, asset_id=asset.id, exc=exc, detail={"upload_resource_id": payload.upload_resource_id, "module": module})
|
||||
raise
|
||||
|
||||
|
||||
@@ -598,10 +706,19 @@ async def soft_delete_asset(db: AsyncSession, *, user_id: str, asset_id: str, li
|
||||
asset.deleted_at = now
|
||||
asset.status = PrivatePortraitAssetStatus.LOCAL_DELETED.value
|
||||
asset.remote_delete_status = PrivatePortraitRemoteDeleteStatus.PENDING.value
|
||||
module = private_portrait_upload_module(asset.library_type)
|
||||
upload_release = await release_upload_resources_by_source(
|
||||
db,
|
||||
source_model=UploadResourceSourceModelEnum.PRIVATE_PORTRAIT_ASSET.value,
|
||||
source_ids=[asset.id],
|
||||
module=module,
|
||||
)
|
||||
setattr(asset, "_pending_upload_resource_ids", list(upload_release.get("released_resource_ids") or []))
|
||||
setattr(asset, "_upload_resource_release", upload_release)
|
||||
await refresh_project_counters(db, [asset.project_id])
|
||||
await db.flush()
|
||||
await db.refresh(asset)
|
||||
log_operation_event(domain=DOMAIN, event_type=PrivatePortraitEventType.ASSET_DELETE_LOCAL.value, event_status=PrivatePortraitEventStatus.SUCCESS.value, source=PrivatePortraitEventSource.API.value, user_id=user_id, project_id=asset.project_id, asset_id=asset.id, detail={"remote_asset_id": asset.remote_asset_id, "remote_project_name": asset.remote_project_name, "library_type": asset.library_type, "asset_type": asset.asset_type})
|
||||
log_operation_event(domain=DOMAIN, event_type=PrivatePortraitEventType.ASSET_DELETE_LOCAL.value, event_status=PrivatePortraitEventStatus.SUCCESS.value, source=PrivatePortraitEventSource.API.value, user_id=user_id, project_id=asset.project_id, asset_id=asset.id, detail={"remote_asset_id": asset.remote_asset_id, "remote_project_name": asset.remote_project_name, "library_type": asset.library_type, "asset_type": asset.asset_type, "upload_resource_release": {k: v for k, v in getattr(asset, "_upload_resource_release", {}).items() if k != "released_resource_ids"}, "pending_upload_resource_count": len(getattr(asset, "_pending_upload_resource_ids", []))})
|
||||
return asset
|
||||
|
||||
|
||||
@@ -617,7 +734,7 @@ async def delete_asset_remote(db: AsyncSession, *, asset_id: str) -> None:
|
||||
log_operation_event(domain=DOMAIN, event_type=PrivatePortraitEventType.ASSET_DELETE_REMOTE_SUCCESS.value, event_status=PrivatePortraitEventStatus.SKIPPED.value, source=PrivatePortraitEventSource.CELERY.value, user_id=asset.user_id, project_id=asset.project_id, asset_id=asset.id, message="远程删除跳过:素材没有 remote_asset_id")
|
||||
return
|
||||
now = datetime.now(timezone.utc)
|
||||
log_operation_event(domain=DOMAIN, event_type=PrivatePortraitEventType.ASSET_DELETE_REMOTE_START.value, event_status=PrivatePortraitEventStatus.PENDING.value, source=PrivatePortraitEventSource.CELERY.value, user_id=asset.user_id, project_id=asset.project_id, asset_id=asset.id, detail={"remote_asset_id": asset.remote_asset_id, "remote_project_name": asset.remote_project_name, "library_type": asset.library_type, "asset_type": asset.asset_type})
|
||||
log_operation_event(domain=DOMAIN, event_type=PrivatePortraitEventType.ASSET_DELETE_REMOTE_START.value, event_status=PrivatePortraitEventStatus.PENDING.value, source=PrivatePortraitEventSource.CELERY.value, user_id=asset.user_id, project_id=asset.project_id, asset_id=asset.id, detail={"remote_asset_id": asset.remote_asset_id, "remote_project_name": asset.remote_project_name, "library_type": asset.library_type, "asset_type": asset.asset_type, "upload_resource_release": {k: v for k, v in getattr(asset, "_upload_resource_release", {}).items() if k != "released_resource_ids"}, "pending_upload_resource_count": len(getattr(asset, "_pending_upload_resource_ids", []))})
|
||||
try:
|
||||
await ArkPrivateAssetClient(for_celery=True).delete_asset(project_name=asset.remote_project_name, asset_id=asset.remote_asset_id)
|
||||
asset.status = PrivatePortraitAssetStatus.REMOTE_DELETED.value
|
||||
|
||||
@@ -19,9 +19,12 @@ from app.enums.private_portrait import (
|
||||
PrivatePortraitProjectStatus,
|
||||
PrivatePortraitRemoteDeleteStatus,
|
||||
)
|
||||
from app.enums.upload_resource import UploadResourceSourceModelEnum
|
||||
from app.models.private_portrait import PrivatePortraitAsset, PrivatePortraitAssetGroup, PrivatePortraitProject
|
||||
from app.schemas.private_portrait import PrivatePortraitProjectCreate, PrivatePortraitProjectOut, PrivatePortraitProjectUpdate
|
||||
from app.services.operation_log_service import log_operation_event
|
||||
from app.services.private_portrait.upload_service import private_portrait_upload_module
|
||||
from app.services.upload_resource import release_upload_resources_by_source
|
||||
from app.utils.id_gen import generate_id
|
||||
|
||||
DOMAIN = "private_portrait"
|
||||
@@ -264,8 +267,29 @@ async def soft_delete_project(
|
||||
) -> PrivatePortraitProject:
|
||||
project = await get_user_project(db, user_id=user_id, project_id=project_id, library_type=library_type)
|
||||
now = datetime.now(timezone.utc)
|
||||
asset_id_rows = await db.execute(
|
||||
select(PrivatePortraitAsset.id)
|
||||
.where(PrivatePortraitAsset.project_id == project_id, PrivatePortraitAsset.deleted_at.is_(None))
|
||||
)
|
||||
asset_ids = [row[0] for row in asset_id_rows.all()]
|
||||
module = private_portrait_upload_module(project.library_type)
|
||||
|
||||
project.deleted_at = now
|
||||
project.status = PrivatePortraitProjectStatus.DELETED.value
|
||||
upload_release = await release_upload_resources_by_source(
|
||||
db,
|
||||
source_model=UploadResourceSourceModelEnum.PRIVATE_PORTRAIT_ASSET.value,
|
||||
source_ids=asset_ids,
|
||||
module=module,
|
||||
) if asset_ids else {
|
||||
"matched": 0,
|
||||
"released": 0,
|
||||
"already_deleted": 0,
|
||||
"released_resource_ids": [],
|
||||
}
|
||||
setattr(project, "_pending_upload_resource_ids", list(upload_release.get("released_resource_ids") or []))
|
||||
setattr(project, "_upload_resource_release", upload_release)
|
||||
|
||||
await db.execute(
|
||||
update(PrivatePortraitAsset)
|
||||
.where(PrivatePortraitAsset.project_id == project_id, PrivatePortraitAsset.deleted_at.is_(None))
|
||||
@@ -285,6 +309,6 @@ async def soft_delete_project(
|
||||
user_id=user_id,
|
||||
project_id=project.id,
|
||||
message="本地软删私域人像素材项目",
|
||||
detail={"library_type": project.library_type, "remote_project_name": project.remote_project_name},
|
||||
detail={"library_type": project.library_type, "remote_project_name": project.remote_project_name, "asset_count": len(asset_ids), "upload_resource_release": {k: v for k, v in upload_release.items() if k != "released_resource_ids"}, "pending_upload_resource_count": len(getattr(project, "_pending_upload_resource_ids", []))},
|
||||
)
|
||||
return project
|
||||
|
||||
@@ -0,0 +1,112 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from fastapi import HTTPException, UploadFile
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.enums.private_portrait import (
|
||||
PrivatePortraitEventSource,
|
||||
PrivatePortraitEventStatus,
|
||||
PrivatePortraitEventType,
|
||||
PrivatePortraitLibraryType,
|
||||
)
|
||||
from app.enums.upload_resource import UploadResourceModuleEnum, UploadResourceTypeEnum
|
||||
from app.models.user import User
|
||||
from app.schemas.private_portrait import PrivatePortraitUploadOut
|
||||
from app.services.operation_log_service import log_operation_error, log_operation_event
|
||||
from app.services.upload_resource import upload_reference_file
|
||||
|
||||
DOMAIN = "private_portrait"
|
||||
PRIVATE_PORTRAIT_IMAGE_MAX_BYTES = 10 * 1024 * 1024
|
||||
PRIVATE_PORTRAIT_VIDEO_MAX_BYTES = 100 * 1024 * 1024
|
||||
|
||||
|
||||
def private_portrait_upload_module(library_type: str) -> str:
|
||||
if library_type == PrivatePortraitLibraryType.REAL_PERSON.value:
|
||||
return UploadResourceModuleEnum.PRIVATE_PORTRAIT_REAL.value
|
||||
if library_type == PrivatePortraitLibraryType.AIGC_VIRTUAL.value:
|
||||
return UploadResourceModuleEnum.PRIVATE_PORTRAIT_VIRTUAL.value
|
||||
raise HTTPException(status_code=400, detail="素材库类型不支持")
|
||||
|
||||
|
||||
async def upload_private_portrait_asset_file(
|
||||
db: AsyncSession,
|
||||
*,
|
||||
file: UploadFile,
|
||||
current_user: User,
|
||||
library_type: str,
|
||||
resource_type: str,
|
||||
duration_seconds: float | None = None,
|
||||
) -> PrivatePortraitUploadOut:
|
||||
if resource_type not in {UploadResourceTypeEnum.IMAGE.value, UploadResourceTypeEnum.VIDEO.value}:
|
||||
raise HTTPException(status_code=400, detail="私域素材上传仅支持图片或视频")
|
||||
|
||||
module = private_portrait_upload_module(library_type)
|
||||
max_bytes = PRIVATE_PORTRAIT_VIDEO_MAX_BYTES if resource_type == UploadResourceTypeEnum.VIDEO.value else PRIVATE_PORTRAIT_IMAGE_MAX_BYTES
|
||||
|
||||
log_operation_event(
|
||||
domain=DOMAIN,
|
||||
event_type=PrivatePortraitEventType.ASSET_UPLOAD_START.value,
|
||||
event_status=PrivatePortraitEventStatus.PENDING.value,
|
||||
source=PrivatePortraitEventSource.API.value,
|
||||
user_id=current_user.id,
|
||||
detail={
|
||||
"module": module,
|
||||
"library_type": library_type,
|
||||
"resource_type": resource_type,
|
||||
"filename": file.filename,
|
||||
"content_type": file.content_type,
|
||||
"duration_seconds": duration_seconds,
|
||||
},
|
||||
)
|
||||
try:
|
||||
result = await upload_reference_file(
|
||||
db,
|
||||
file=file,
|
||||
current_user=current_user,
|
||||
module=module,
|
||||
resource_type=resource_type,
|
||||
gen_type="private_portrait",
|
||||
duration_seconds=duration_seconds,
|
||||
max_bytes=max_bytes,
|
||||
)
|
||||
out = PrivatePortraitUploadOut(
|
||||
url=result.url,
|
||||
filename=result.filename,
|
||||
type=result.resource_type,
|
||||
module=result.module,
|
||||
resource_id=result.resource_id,
|
||||
file_size_bytes=result.file_size_bytes,
|
||||
duration_seconds=result.duration_seconds,
|
||||
)
|
||||
log_operation_event(
|
||||
domain=DOMAIN,
|
||||
event_type=PrivatePortraitEventType.ASSET_UPLOAD_SUCCESS.value,
|
||||
event_status=PrivatePortraitEventStatus.SUCCESS.value,
|
||||
source=PrivatePortraitEventSource.API.value,
|
||||
user_id=current_user.id,
|
||||
detail={
|
||||
"module": module,
|
||||
"library_type": library_type,
|
||||
"resource_type": resource_type,
|
||||
"resource_id": result.resource_id,
|
||||
"url": result.url,
|
||||
"file_size_bytes": result.file_size_bytes,
|
||||
"duration_seconds": result.duration_seconds,
|
||||
},
|
||||
)
|
||||
return out
|
||||
except Exception as exc: # noqa: BLE001
|
||||
log_operation_error(
|
||||
domain=DOMAIN,
|
||||
event_type=PrivatePortraitEventType.ASSET_UPLOAD_FAILED.value,
|
||||
source=PrivatePortraitEventSource.API.value,
|
||||
user_id=current_user.id,
|
||||
exc=exc,
|
||||
detail={
|
||||
"module": module,
|
||||
"library_type": library_type,
|
||||
"resource_type": resource_type,
|
||||
"filename": file.filename,
|
||||
},
|
||||
)
|
||||
raise
|
||||
@@ -2,6 +2,7 @@ from __future__ import annotations
|
||||
|
||||
import uuid
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from fastapi import HTTPException
|
||||
@@ -29,6 +30,7 @@ from app.schemas.shot_replicate import (
|
||||
ShotSegmentDeleteOut,
|
||||
ShotSegmentDetailOut,
|
||||
ShotSegmentListOut,
|
||||
ShotSegmentSplitRetryOut,
|
||||
ShotReanalyzeOut,
|
||||
ShotSegmentOut,
|
||||
ShotSplitByAIOut,
|
||||
@@ -500,6 +502,97 @@ async def create_segments_by_ai(
|
||||
)
|
||||
|
||||
|
||||
async def prepare_retry_split_segment(
|
||||
db: AsyncSession,
|
||||
*,
|
||||
current_user: User,
|
||||
segment_id: str,
|
||||
force: bool = False,
|
||||
reason: str | None = None,
|
||||
) -> ShotSegmentSplitRetryOut:
|
||||
"""重置失败切片片段,commit 成功后由 API 投递现有 split_one_segment 任务。"""
|
||||
segment = await get_segment_for_user(db, segment_id=segment_id, user=current_user, for_update=True)
|
||||
task_set = await get_task_set_for_user(db, task_set_id=segment.task_set_id, user=current_user, for_update=True)
|
||||
from_split_status = segment.split_status
|
||||
allowed = {ShotSplitStatusEnum.FAILED.value, ShotSplitStatusEnum.RETRY_WAITING.value}
|
||||
|
||||
reject_reason: str | None = None
|
||||
if task_set.deleted_at is not None or task_set.status == ShotTaskSetStatusEnum.DELETED.value:
|
||||
reject_reason = "拆镜总任务集已删除,不能重试切片"
|
||||
elif segment.deleted_at is not None:
|
||||
reject_reason = "拆镜片段已删除,不能重试切片"
|
||||
elif from_split_status == ShotSplitStatusEnum.PROCESSING.value:
|
||||
reject_reason = "拆镜片段正在切片处理中,不能重复投递"
|
||||
elif from_split_status == ShotSplitStatusEnum.COMPLETED.value:
|
||||
reject_reason = "拆镜片段已切片完成,不支持重切,避免旧切片资源覆盖"
|
||||
elif from_split_status not in allowed and not force:
|
||||
reject_reason = "仅允许失败或等待重试的切片片段重新投递"
|
||||
elif not task_set.video_path:
|
||||
reject_reason = "原视频本地路径为空,不能重试切片"
|
||||
elif not Path(str(task_set.video_path)).exists():
|
||||
reject_reason = "原视频本地文件不存在,不能重试切片"
|
||||
|
||||
if reject_reason:
|
||||
log_module_event_file(
|
||||
module=MODULE,
|
||||
event_type=ShotReplicateLogEventEnum.SEGMENT_SPLIT_RETRY_REJECTED.value,
|
||||
project_id=task_set.id,
|
||||
step_id=segment.id,
|
||||
status="rejected",
|
||||
message=reject_reason,
|
||||
detail={
|
||||
"segment_id": segment.id,
|
||||
"task_set_id": task_set.id,
|
||||
"from_split_status": from_split_status,
|
||||
"force": force,
|
||||
"reason": reason,
|
||||
},
|
||||
)
|
||||
raise HTTPException(status_code=400, detail=reject_reason)
|
||||
|
||||
now = _now()
|
||||
segment.split_status = ShotSplitStatusEnum.PENDING.value
|
||||
segment.split_enqueued_at = now
|
||||
segment.split_started_at = None
|
||||
segment.split_lease_until = None
|
||||
segment.split_next_retry_at = None
|
||||
segment.split_retry_count = 0
|
||||
segment.split_last_error = None
|
||||
segment.split_celery_task_id = f"shot-split:{uuid.uuid4().hex}"
|
||||
task_set.split_error_message = None
|
||||
|
||||
await refresh_task_set_split_summary(db, task_set.id)
|
||||
await db.flush()
|
||||
|
||||
log_module_event_file(
|
||||
module=MODULE,
|
||||
event_type=ShotReplicateLogEventEnum.SEGMENT_SPLIT_RETRY_RECEIVED.value,
|
||||
project_id=task_set.id,
|
||||
step_id=segment.id,
|
||||
status="pending",
|
||||
message="拆镜片段切片失败重试已重置,等待投递 Celery",
|
||||
detail={
|
||||
"segment_id": segment.id,
|
||||
"task_set_id": task_set.id,
|
||||
"from_split_status": from_split_status,
|
||||
"to_split_status": segment.split_status,
|
||||
"force": force,
|
||||
"reason": reason,
|
||||
"source_path": task_set.video_path,
|
||||
"celery_task_name": "shot_replicate.split_one_segment",
|
||||
"queue": "gen_result_download",
|
||||
},
|
||||
)
|
||||
|
||||
return ShotSegmentSplitRetryOut(
|
||||
message="切片重试已提交,正在重新切割视频片段",
|
||||
task_set_id=task_set.id,
|
||||
segment_id=segment.id,
|
||||
split_status=segment.split_status,
|
||||
celery_task_name="shot_replicate.split_one_segment",
|
||||
)
|
||||
|
||||
|
||||
async def create_custom_segment(
|
||||
db: AsyncSession,
|
||||
*,
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Any, Iterable
|
||||
|
||||
from sqlalchemy import select
|
||||
|
||||
@@ -15,7 +15,7 @@ from app.enums.upload_resource import UploadResourceModuleEnum, UploadResourceTy
|
||||
COMMON_IMAGE_RE = re.compile(r"^images/(?P<year>\d{4})/(?P<month>\d{2})/(?P<day>\d{2})/video_img_(?P<user_id>[^_]+)_(?P<ymd>\d{8})_(?P<hms>\d{6})_(?P<rand>[0-9a-fA-F]{8})\.(?P<ext>[^/]+)$")
|
||||
COMMON_VIDEO_RE = re.compile(r"^videos/(?P<year>\d{4})/(?P<month>\d{2})/(?P<day>\d{2})/video_ref_(?P<user_id>[^_]+)_(?P<ymd>\d{8})_(?P<hms>\d{6})_(?P<rand>[0-9a-fA-F]{8})\.(?P<ext>[^/]+)$")
|
||||
COMMON_AUDIO_RE = re.compile(r"^audios/(?P<year>\d{4})/(?P<month>\d{2})/(?P<day>\d{2})/audio_ref_(?P<user_id>[^_]+)_(?P<ymd>\d{8})_(?P<hms>\d{6})_(?P<rand>[0-9a-fA-F]{8})\.(?P<ext>[^/]+)$")
|
||||
MODULE_RE = re.compile(r"^(?P<module>hot_opening_replicate|shot_replicate)/(?P<kind>images|videos)/(?P<year>\d{4})/(?P<month>\d{2})/(?P<day>\d{2})/(?P<prefix>video_img|video_ref)_(?P<user_id>[^_]+)_(?P<ymd>\d{8})_(?P<hms>\d{6})_(?P<rand>[0-9a-fA-F]{8})\.(?P<ext>[^/]+)$")
|
||||
MODULE_RE = re.compile(r"^(?P<module>hot_opening_replicate|shot_replicate|private_portrait_real|private_portrait_virtual)/(?P<kind>images|videos)/(?P<year>\d{4})/(?P<month>\d{2})/(?P<day>\d{2})/(?P<prefix>video_img|video_ref)_(?P<user_id>[^_]+)_(?P<ymd>\d{8})_(?P<hms>\d{6})_(?P<rand>[0-9a-fA-F]{8})\.(?P<ext>[^/]+)$")
|
||||
SHOT_SEGMENT_RE = re.compile(r"^shot_segments/(?P<year>\d{4})/(?P<month>\d{2})/(?P<day>\d{2})/(?P<segment_id>[^/]+)\.mp4$")
|
||||
LEGACY_GEN_RE = re.compile(r"^(?P<kind>images|videos)/gen_(?P<user_id>[^_]+)_(?P<rand>[0-9a-zA-Z]+)\.(?P<ext>[^/]+)$")
|
||||
ADMIN_UPLOAD_RE = re.compile(r"^admin_uploads/(?P<scene>system_logo|system_pdf|open_type_thumb|common)/(?P<kind>images|videos|audios|files)/(?P<year>\d{4})/(?P<month>\d{2})/(?P<day>\d{2})/(?P<prefix>admin_img|admin_video|admin_audio|admin_pdf|admin_file)_(?P<user_id>[^_]+)_(?P<ymd>\d{8})_(?P<hms>\d{6})_(?P<rand>[0-9a-fA-F]{8})\.(?P<ext>[^/]+)$")
|
||||
|
||||
Reference in New Issue
Block a user