from __future__ import annotations from urllib.parse import unquote from fastapi import APIRouter, Depends, HTTPException, Query, Request from fastapi.responses import RedirectResponse from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from app.dependencies import get_current_user, get_db from app.enums.private_portrait import ( PrivatePortraitEventSource, PrivatePortraitEventStatus, PrivatePortraitEventType, PrivatePortraitRemoteDeleteStatus, ) from app.models.private_portrait import PrivatePortraitAsset, PrivatePortraitProject from app.models.user import User from app.schemas.private_portrait import ( PrivatePortraitAssetCreate, PrivatePortraitAssetListOut, PrivatePortraitDeleteOut, PrivatePortraitConfigOut, PrivatePortraitProjectCreate, PrivatePortraitProjectListOut, PrivatePortraitProjectOut, PrivatePortraitProjectUpdate, PrivatePortraitSelectableAssetListOut, PrivatePortraitValidateSessionCreate, PrivatePortraitValidateSessionOut, ) from app.services.operation_log_service import log_operation_error, log_operation_event from app.services.private_portrait.asset_service import ( DOMAIN, asset_to_out, create_asset, create_validate_session, get_user_private_portrait_config, get_validate_session, handle_validate_callback, list_assets, list_selectable_assets, soft_delete_asset, sync_asset_status, validate_session_to_out, ) from app.services.private_portrait.project_service import ( create_project, get_user_project, list_projects, project_to_out, refresh_project_counters, soft_delete_project, update_project, ) router = APIRouter(tags=["private-portrait"]) 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}, ) 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}, ) @router.get("/private-portrait/config", response_model=PrivatePortraitConfigOut) 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) @router.post("/private-portrait/projects", response_model=PrivatePortraitProjectOut) async def create_private_portrait_project(payload: PrivatePortraitProjectCreate, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db)): project = await create_project(db, user_id=current_user.id, payload=payload) out = project_to_out(project) await db.commit() return out @router.get("/private-portrait/projects", response_model=PrivatePortraitProjectListOut) async def list_private_portrait_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), ): items, total = await list_projects(db, user_id=current_user.id, page=page, page_size=page_size, keyword=keyword, status=status) await refresh_project_counters(db, [item.id for item in items]) await db.commit() items, total = await list_projects(db, user_id=current_user.id, page=page, page_size=page_size, keyword=keyword, status=status) 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) 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)) @router.put("/private-portrait/projects/{project_id}", response_model=PrivatePortraitProjectOut) 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_project(db, user_id=current_user.id, project_id=project_id, payload=payload) out = project_to_out(project) await db.commit() return out @router.delete("/private-portrait/projects/{project_id}", response_model=PrivatePortraitDeleteOut) 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) project_id_snapshot = project.id await db.commit() try: from app.tasks.private_portrait_asset_tasks import delete_private_portrait_project_remote delete_private_portrait_project_remote.delay(project_id_snapshot) _log_task_dispatch_success(task_name="private_portrait.delete_project_remote", user_id=current_user.id, project_id=project_id_snapshot) except Exception as exc: _log_task_dispatch_failed(task_name="private_portrait.delete_project_remote", user_id=current_user.id, project_id=project_id_snapshot, exc=exc) return PrivatePortraitDeleteOut(success=True, remote_delete_status=PrivatePortraitRemoteDeleteStatus.PENDING.value) @router.post("/private-portrait/projects/{project_id}/validate-sessions", response_model=PrivatePortraitValidateSessionOut) 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_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) await db.commit() return out @router.get("/private-portrait/validate-sessions/{session_id}", response_model=PrivatePortraitValidateSessionOut) 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") 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) params.pop("redirect_url", None) session = await handle_validate_callback(db, session_id=session_id, query_params=params) redirect_session_id = session.id redirect_status = session.status redirect_result_code = session.result_code or "" response = {"session_id": session.id, "status": session.status, "resultCode": session.result_code, "remote_group_id": session.remote_group_id} await db.commit() if redirect_url: sep = "&" if "?" in redirect_url else "?" url = f"{unquote(redirect_url)}{sep}session_id={redirect_session_id}&status={redirect_status}&resultCode={redirect_result_code}" return RedirectResponse(url=url) return response @router.post("/private-portrait/projects/{project_id}/assets") 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_asset(db, user_id=current_user.id, project_id=project_id, payload=payload) asset_id_snapshot = asset.id project_id_snapshot = asset.project_id out = asset_to_out(asset) await db.commit() try: from app.tasks.private_portrait_asset_tasks import poll_private_portrait_asset_status poll_private_portrait_asset_status.delay(asset_id_snapshot) _log_task_dispatch_success(task_name="private_portrait.poll_asset_status", user_id=current_user.id, project_id=project_id_snapshot, asset_id=asset_id_snapshot) except Exception as exc: _log_task_dispatch_failed(task_name="private_portrait.poll_asset_status", user_id=current_user.id, project_id=project_id_snapshot, asset_id=asset_id_snapshot, exc=exc) return out @router.get("/private-portrait/projects/{project_id}/assets", response_model=PrivatePortraitAssetListOut) 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), 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) 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}") 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).limit(1))).scalar_one_or_none() if not asset: raise HTTPException(status_code=404, detail="真人素材不存在") project = (await db.execute(select(PrivatePortraitProject).where(PrivatePortraitProject.id == asset.project_id).limit(1))).scalar_one_or_none() return asset_to_out(asset, project_name=project.name if project else None) @router.post("/private-portrait/assets/{asset_id}/sync") 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) out = asset_to_out(asset) await db.commit() return out @router.delete("/private-portrait/assets/{asset_id}", response_model=PrivatePortraitDeleteOut) 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) asset_id_snapshot = asset.id project_id_snapshot = asset.project_id await db.commit() try: from app.tasks.private_portrait_asset_tasks import delete_private_portrait_asset_remote delete_private_portrait_asset_remote.delay(asset_id_snapshot) _log_task_dispatch_success(task_name="private_portrait.delete_asset_remote", user_id=current_user.id, project_id=project_id_snapshot, asset_id=asset_id_snapshot) except Exception as exc: _log_task_dispatch_failed(task_name="private_portrait.delete_asset_remote", user_id=current_user.id, project_id=project_id_snapshot, asset_id=asset_id_snapshot, exc=exc) return PrivatePortraitDeleteOut(success=True, remote_delete_status=PrivatePortraitRemoteDeleteStatus.PENDING.value) @router.get("/private-portrait/selectable-assets", response_model=PrivatePortraitSelectableAssetListOut) 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), 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) return PrivatePortraitSelectableAssetListOut(items=items, total=total, page=page, page_size=page_size)