diff --git a/apps/worker/worker_app/tasks/asset_quality_scoring_task.py b/apps/worker/worker_app/tasks/asset_quality_scoring_task.py index 083b58c49..83a265acd 100644 --- a/apps/worker/worker_app/tasks/asset_quality_scoring_task.py +++ b/apps/worker/worker_app/tasks/asset_quality_scoring_task.py @@ -18,6 +18,7 @@ from worker_app.celery_app import celery_app from worker_app.db import SessionLocal from packages.adapters.sqlalchemy_impl.asset_repository import SQLAlchemyAssetRepository +from packages.domain.classification import ClassificationStatus from packages.shared.storage import get_shared_storage_service logger = get_task_logger(__name__) @@ -96,7 +97,7 @@ def calculate_asset_quality_task(self, asset_id: str) -> dict: confidence = 1.0 existing_meta["classification"] = classification existing_meta["classification_confidence"] = confidence - asset.classification_status = "completed" + asset.classification_status = ClassificationStatus.COMPLETED asset.metadata = existing_meta logger.info( "[quality_score] asset=%s 自动分类完成: category=%s confidence=%.2f", @@ -110,6 +111,8 @@ def calculate_asset_quality_task(self, asset_id: str) -> dict: asset_id, cls_err, ) + # 分类失败显式标记 FAILED,避免停留在 PENDING 被反复重试 + asset.classification_status = ClassificationStatus.FAILED asset_repo.update(asset) db.commit() diff --git a/packages/adapters/sqlalchemy_impl/asset_repository.py b/packages/adapters/sqlalchemy_impl/asset_repository.py index f316fa543..d1cab4535 100755 --- a/packages/adapters/sqlalchemy_impl/asset_repository.py +++ b/packages/adapters/sqlalchemy_impl/asset_repository.py @@ -128,8 +128,12 @@ class SQLAlchemyAssetRepository: height=asset.height, fps=asset.fps, codec=asset.codec, - status=asset.status.value, - classification_status=asset.classification_status.value, + status=(asset.status.value if hasattr(asset.status, "value") else str(asset.status)), + classification_status=( + asset.classification_status.value + if hasattr(asset.classification_status, "value") + else str(asset.classification_status) + ), classification_result=(json.dumps(asset.metadata) if asset.metadata else None), quality_score=asset.quality_score, uploaded_by_user_id=asset.uploaded_by_user_id or "system", @@ -142,7 +146,7 @@ class SQLAlchemyAssetRepository: self.session.flush() self._sync_asset_tags(asset.id, asset.tag_ids) # Issue #1776: 自动维护素材库计数(同事务内原子更新) - if asset.library_id and asset.status.value != "deleted": + if asset.library_id and (getattr(asset.status, "value", str(asset.status)) != "deleted"): from sqlalchemy import func self.session.query(AssetLibraryModel).filter(AssetLibraryModel.id == asset.library_id).update( @@ -168,8 +172,12 @@ class SQLAlchemyAssetRepository: model.height = asset.height model.fps = asset.fps model.codec = asset.codec - model.status = asset.status.value - model.classification_status = asset.classification_status.value + model.status = asset.status.value if hasattr(asset.status, "value") else str(asset.status) + model.classification_status = ( + asset.classification_status.value + if hasattr(asset.classification_status, "value") + else str(asset.classification_status) + ) model.classification_result = json.dumps(asset.metadata) if asset.metadata else None model.quality_score = asset.quality_score model.uploaded_by_user_id = asset.uploaded_by_user_id or model.uploaded_by_user_id