diff --git a/apps/api/app/services/lipsync_service.py b/apps/api/app/services/lipsync_service.py index 442b2c970..283c6f032 100644 --- a/apps/api/app/services/lipsync_service.py +++ b/apps/api/app/services/lipsync_service.py @@ -198,6 +198,12 @@ class LipsyncService: self.db.add(job) self.db.flush() + # ⚠️ 必须先 commit 再发 Celery 任务,避免事务竞态: + # worker 是独立进程+独立DB连接,任务被消费(<4ms)时若本事务还未提交, + # worker 查询 job 会返回 None → 静默 return 不重试,job 永远卡在 tts_processing。 + self.db.commit() + self.db.refresh(job) + if is_tts_mode: # 2a. TTS 模式:dispatch Celery 异步任务处理 TTS 合成 + MediaKit 提交 try: @@ -223,6 +229,7 @@ class LipsyncService: job.error_message = f"Celery 任务投递失败: {exc}" job.error_code = "AsyncDispatchFailed" job.updated_at = datetime.now(timezone.utc) + self.db.commit() # 投递失败也要落库失败状态 else: # 2b. 直接音频模式:同步签名并提交 MediaKit video_url = self._sign_media_url(video_url) @@ -240,15 +247,15 @@ class LipsyncService: job.mediakit_task_id = result["task_id"] job.status = "submitted" job.submitted_at = datetime.now(timezone.utc) + self.db.commit() # submitted 状态落库 except MediaKitError as exc: job.status = "failed" job.error_message = str(exc) job.error_code = exc.code logger.error("提交对口型任务失败: %s", exc) + self.db.commit() raise - self.db.commit() - self.db.refresh(job) return job # ── 查询任务 ────────────────────────────────────────────────────────── diff --git a/apps/api/app/tasks/lipsync_tts.py b/apps/api/app/tasks/lipsync_tts.py index 1b3ef0dfc..7651911f5 100644 --- a/apps/api/app/tasks/lipsync_tts.py +++ b/apps/api/app/tasks/lipsync_tts.py @@ -211,8 +211,11 @@ def _estimate_sentence_timings_by_chars(sentences: list[str], total_duration: fl @shared_task( bind=True, name="lipsync_tts.synthesize_and_submit", - max_retries=2, + max_retries=5, # 事务竞态重试3次(job not found)+ TTS偶发错误2次 default_retry_delay=30, + autoretry_for=(OSError, ConnectionError), # 网络/连接错误自动重试 + retry_backoff=True, + retry_backoff_max=30, soft_time_limit=180, time_limit=200, ) @@ -259,7 +262,28 @@ def tts_synthesize_and_submit( ) if job is None: - logger.error("[lipsync_tts] Job not found: job_id=%s", job_id) + # 事务竞态防御:API 在 commit 前投递了任务,worker 消费时事务尚未提交。 + # Celery 内置 autoretry_for 不支持"业务条件重试",这里手动 retry 3 次, + # 间隔递增(1s/3s/7s),让 API 事务有时间提交。 + # max_retries 由 self.request(retries) 维护;默认 self.max_retries=3 由装饰器 soft_time_limit 下方指定。 + retries = getattr(self.request, "retries", 0) + max_retries = 3 + if retries < max_retries: + backoff = (2 ** retries) + (retries * 1) # 1s, 3s, 7s + logger.warning( + "[lipsync_tts] Job not found yet (retry %d/%d, backoff %ds): job_id=%s", + retries + 1, + max_retries, + backoff, + job_id, + ) + self.db.close() + raise self.retry(countdown=backoff, max_retries=max_retries) + logger.error( + "[lipsync_tts] Job not found after %d retries, giving up: job_id=%s", + max_retries, + job_id, + ) return # 已取消的任务不再处理 diff --git a/tests/unit/test_lipsync_speed_optimization.py b/tests/unit/test_lipsync_speed_optimization.py index 5963b8a51..13d9e4220 100644 --- a/tests/unit/test_lipsync_speed_optimization.py +++ b/tests/unit/test_lipsync_speed_optimization.py @@ -288,3 +288,44 @@ class TestCancelJobTtsProcessing: result = svc.cancel_job("job-1", "user-1") assert result.status == "cancelled" + + +class TestCreateJobCommitOrder: + """验证事务顺序修复:create_job 必须先 commit 再发 Celery 任务,避免 worker 消费时 job 不可见。""" + + def test_commit_called_before_apply_async_in_tts_mode(self): + """TTS 模式:db.commit() 必须在 apply_async() 之前调用,防止 worker 查不到 job 永远卡在 tts_processing。""" + svc, client, cosy = _make_service_with_mocks() + call_order: list[str] = [] + + def track_commit(): + call_order.append("commit") + def track_apply_async(*args, **kwargs): + call_order.append("apply_async") + + svc.db.commit.side_effect = track_commit + + with patch("app.services.lipsync_service.tts_synthesize_and_submit") as mock_task: + mock_task.apply_async = MagicMock(side_effect=track_apply_async) + svc.create_job( + user_id="user-1", + video_url="https://example.com/video.mp4", + voice_id="v-1", + script_text="测试", + ) + + # 至少有一次 commit 在 apply_async 之前 + assert "commit" in call_order, "db.commit 必须被调用" + assert "apply_async" in call_order, "apply_async 必须被调用" + assert call_order.index("commit") < call_order.index("apply_async"), ( + f"事务顺序错误:commit 必须在 apply_async 之前,实际顺序 {call_order}" + ) + + def test_job_not_found_retry_mechanism_exists(self): + """worker 侧 job not found 必须有重试机制(self.retry),而不是静默 return。""" + import inspect + from app.tasks.lipsync_tts import tts_synthesize_and_submit + source = inspect.getsource(tts_synthesize_and_submit.run) + assert "self.retry" in source or "retry" in source, ( + "tts_synthesize_and_submit 在 job not found 时必须重试,防止静默失败" + )