From 642692684494f14718ccc2c50ed835559e54811a Mon Sep 17 00:00:00 2001 From: saas-backend-agent Date: Thu, 3 Sep 2026 17:05:12 +0800 Subject: [PATCH] =?UTF-8?q?fix(worker):=20Worker=20=E5=BC=82=E5=B8=B8?= =?UTF-8?q?=E6=97=B6=E5=B0=86=E5=8D=A0=E4=BD=8D=20Asset=20=E6=A0=87?= =?UTF-8?q?=E8=AE=B0=E4=B8=BA=20ERROR=EF=BC=8C=E9=81=BF=E5=85=8D=E6=B0=B8?= =?UTF-8?q?=E8=BF=9C=E5=8D=A1=E5=9C=A8=20PROCESSING?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ingest_asset 的 except 块原来只更新 job status 为 FAILED,没有处理 上传时预先创建的 PROCESSING 状态 Asset。如果 Worker 在处理过程中 崩溃(ffmpeg 失败、OSS 超时等),Asset 会永远卡在 PROCESSING。 修复:在 except 块中查找对应 storage_key 的 Asset,如果仍处于 PROCESSING 或 UPLOADING 状态,将其标记为 ERROR 并记录错误信息。 asset_repo 查找失败不影响 job status 更新(内层 try/except 隔离)。 --- apps/worker/worker_app/tasks/ingest.py | 24 +++++++++++++++++++++++- 1 file changed, 23 insertions(+), 1 deletion(-) 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() -- 2.54.0