diff --git a/apps/worker/worker_app/tasks/ingest.py b/apps/worker/worker_app/tasks/ingest.py index 7a9f362de..0a7aa4764 100755 --- a/apps/worker/worker_app/tasks/ingest.py +++ b/apps/worker/worker_app/tasks/ingest.py @@ -723,15 +723,37 @@ def ingest_asset(job_id: str) -> dict: db.rollback() logger.error(f"Failed to ingest asset {job_id}: {e}") - # Update job status to FAILED + # Update job status to FAILED and mark pre-created Asset as ERROR try: job_repo = SQLAlchemyIngestJobRepository(db) + asset_repo = SQLAlchemyAssetRepository(db) job = job_repo.get(job_id) if job: job.status = IngestJobStatus.FAILED job.error_message = str(e) job.updated_at = datetime.now(timezone.utc) job_repo.update(job) + + # 将上传时创建的占位 Asset(PROCESSING/UPLOADING)标记为 ERROR, + # 避免素材永远卡在中间状态 + try: + existing = asset_repo.find_by_storage_key(job.storage_key) + if existing and existing.status in ( + AssetStatus.PROCESSING, + AssetStatus.UPLOADING, + ): + existing.status = AssetStatus.ERROR + existing.metadata = {**(existing.metadata or {}), "ingest_error": str(e)} + existing.updated_at = datetime.now(timezone.utc) + asset_repo.update(existing) + logger.info( + "Marked asset as ERROR due to ingest failure: asset_id=%s job_id=%s", + existing.id, + job_id, + ) + except Exception as asset_err: + logger.warning("Failed to mark asset as ERROR: %s", asset_err) + db.commit() except Exception: db.rollback()