Files

205 lines
7.9 KiB
Python

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