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()