from __future__ import annotations from datetime import datetime, timezone from pathlib import Path from typing import Iterable, Any from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from app.enums.upload_resource import UploadResourceEventEnum, UploadResourceFileDeleteStatusEnum from app.models.upload_resource import UploadResource from app.services.upload_resource.log_service import log_upload_resource_event, log_upload_resource_exception from app.services.upload_resource.path_resolver import normalize_storage_path def _clean_ids(values: Iterable[str | None] | None) -> list[str]: if not values: return [] return [str(v).strip() for v in dict.fromkeys(values) if v and str(v).strip()] def _clean_paths(values: Iterable[str | None] | None) -> list[str]: if not values: return [] cleaned: list[str] = [] for value in values: if not value: continue try: cleaned.append(normalize_storage_path(value)) except Exception: cleaned.append(str(value)) return list(dict.fromkeys(cleaned)) def _short_error(exc: BaseException) -> str: text = str(exc) or exc.__class__.__name__ return text[:2000] async def cleanup_upload_resource_files_after_commit( db: AsyncSession, *, resource_ids: Iterable[str | None] | None = None, legacy_paths: Iterable[str | None] | None = None, ) -> dict[str, int]: """API 主事务 commit 成功后清理 UploadResource 真实文件。 这里不负责业务删除事务,不做 rollback;调用方如需持久化清理状态, 应在本函数返回后由 API 层再次 commit。 """ ids = _clean_ids(resource_ids) paths = _clean_paths(legacy_paths) stats = { "matched": 0, "deleted": 0, "missing": 0, "failed": 0, "legacy_deleted": 0, "legacy_missing": 0, "legacy_failed": 0, } if ids: try: result = await db.execute( select(UploadResource) .where( UploadResource.id.in_(ids), UploadResource.deleted_at.is_not(None), UploadResource.physical_deleted_at.is_(None), UploadResource.file_delete_status.in_([ UploadResourceFileDeleteStatusEnum.PENDING_DELETE.value, UploadResourceFileDeleteStatusEnum.DELETE_FAILED.value, ]), ) .with_for_update() ) resources = list(result.scalars().all()) except Exception as exc: # noqa: BLE001 log_upload_resource_exception( event_type=UploadResourceEventEnum.UPLOAD_RESOURCE_CLEANUP_BATCH_FAILED.value, message="UploadResource 真实文件清理批量查询失败", resource_ids=ids, detail={"stage": "query_resources", "resource_ids_count": len(ids)}, exc=exc, ) raise stats["matched"] = len(resources) for resource in resources: now = datetime.now(timezone.utc) try: path = Path(resource.storage_path) if path.exists(): path.unlink() resource.file_delete_status = UploadResourceFileDeleteStatusEnum.DELETED.value stats["deleted"] += 1 event = UploadResourceEventEnum.DELETE_PHYSICAL_SUCCESS.value else: resource.file_delete_status = UploadResourceFileDeleteStatusEnum.MISSING.value stats["missing"] += 1 event = UploadResourceEventEnum.DELETE_PHYSICAL_MISSING.value resource.physical_deleted_at = now resource.file_delete_error = None log_upload_resource_event( event_type=event, module=resource.module, user_id=resource.user_id, resource_id=resource.id, source_model=resource.source_model, source_id=resource.source_id, detail={"storage_path": resource.storage_path, "file_delete_status": resource.file_delete_status}, ) except Exception as exc: # noqa: BLE001 resource.file_delete_status = UploadResourceFileDeleteStatusEnum.DELETE_FAILED.value resource.file_delete_error = _short_error(exc) stats["failed"] += 1 log_upload_resource_event( event_type=UploadResourceEventEnum.DELETE_PHYSICAL_FAILED.value, module=resource.module, user_id=resource.user_id, resource_id=resource.id, source_model=resource.source_model, source_id=resource.source_id, detail={"storage_path": resource.storage_path}, exc=exc, ) if resources: try: await db.flush() except Exception as exc: # noqa: BLE001 log_upload_resource_exception( event_type=UploadResourceEventEnum.UPLOAD_RESOURCE_CLEANUP_BATCH_FAILED.value, message="UploadResource 真实文件清理状态 flush 失败", resource_ids=[resource.id for resource in resources], detail={"stage": "flush_cleanup_status"}, exc=exc, ) raise for raw_path in paths: try: path = Path(raw_path) if path.exists(): path.unlink() stats["legacy_deleted"] += 1 else: stats["legacy_missing"] += 1 except Exception as exc: # noqa: BLE001 stats["legacy_failed"] += 1 log_upload_resource_event( event_type=UploadResourceEventEnum.DELETE_PHYSICAL_FAILED.value, detail={"legacy_path": raw_path}, exc=exc, ) return stats async def cleanup_pending_upload_resource_files(db: AsyncSession, *, limit: int = 500) -> dict[str, int]: limit = max(1, int(limit or 500)) log_upload_resource_event( event_type=UploadResourceEventEnum.CLEANUP_PENDING_START.value, detail={"limit": limit}, ) try: result = await db.execute( select(UploadResource.id) .where( UploadResource.deleted_at.is_not(None), UploadResource.physical_deleted_at.is_(None), UploadResource.file_delete_status.in_([ UploadResourceFileDeleteStatusEnum.PENDING_DELETE.value, UploadResourceFileDeleteStatusEnum.DELETE_FAILED.value, ]), ) .order_by(UploadResource.updated_at.asc()) .limit(limit) ) ids = list(result.scalars().all()) except Exception as exc: # noqa: BLE001 log_upload_resource_exception( event_type=UploadResourceEventEnum.UPLOAD_RESOURCE_CLEANUP_BATCH_FAILED.value, message="UploadResource pending 清理查询失败", detail={"stage": "query_pending_cleanup", "limit": limit}, exc=exc, ) raise try: stats = await cleanup_upload_resource_files_after_commit(db, resource_ids=ids) except Exception as exc: # noqa: BLE001 log_upload_resource_exception( event_type=UploadResourceEventEnum.UPLOAD_RESOURCE_CLEANUP_BATCH_FAILED.value, message="UploadResource pending 真实文件补偿清理失败", resource_ids=ids, detail={"stage": "cleanup_pending", "limit": limit}, exc=exc, ) raise log_upload_resource_event( event_type=UploadResourceEventEnum.CLEANUP_PENDING_FINISHED.value, detail={"ids": len(ids), **stats}, ) return stats