diff --git a/apps/api/app/api/routes/lipsync.py b/apps/api/app/api/routes/lipsync.py index c0d989dea..7fc23a56e 100644 --- a/apps/api/app/api/routes/lipsync.py +++ b/apps/api/app/api/routes/lipsync.py @@ -14,7 +14,6 @@ import logging from app.auth import AuthenticatedUser, get_current_user from app.dependencies import ( - get_cosyvoice_service, get_db_session, get_voice_clone_profile_repository, ) @@ -24,8 +23,6 @@ from app.services.mediakit_client import MediaKitError from fastapi import APIRouter, BackgroundTasks, Depends, HTTPException, Query from sqlalchemy.orm import Session -from packages.application.cosyvoice_service import CosyVoiceError, CosyVoiceService - logger = logging.getLogger(__name__) router = APIRouter() @@ -34,13 +31,11 @@ router = APIRouter() def _get_service( db: Session = Depends(get_db_session), voice_clone_repo=Depends(get_voice_clone_profile_repository), - cosyvoice_service: CosyVoiceService = Depends(get_cosyvoice_service), ) -> LipsyncService: - # voice_clone_repo 用于克隆音色 profile 解析;cosyvoice_service 用于 TTS 直生 - # (TTS 合成、音色解析、错误码归一化都在 LipsyncService 内部完成) + # voice_clone_repo 用于克隆音色 profile 解析 + # TTS 合成已移至 Celery 异步任务,无需同步注入 cosyvoice_service return LipsyncService( db, - cosyvoice_service=cosyvoice_service, voice_clone_repo=voice_clone_repo, ) @@ -57,8 +52,8 @@ def create_lipsync_job( """提交对口型任务. #1809/#1822: 前端传 {video_url, voice_id, script_text, speed?, emotion?}, - 后端内部解析音色、调 TTS 合成音频、转存 OSS,再提交 MediaKit; - 也支持直接传 {video_url, audio_url}。 + 后端创建任务记录(状态 tts_processing),dispatch Celery 异步任务执行 TTS 合成 + MediaKit 提交; + 也支持直接传 {video_url, audio_url}(同步提交 MediaKit)。 """ try: job = svc.create_job( @@ -75,21 +70,13 @@ def create_lipsync_job( except ValueError as exc: # 参数无效(如 voice_id 格式不对、文本过长等) raise HTTPException(status_code=400, detail=str(exc)) from exc - except CosyVoiceError as exc: - # TTS 合成基础设施失败(API/网络/认证) - raise HTTPException( - status_code=502, - detail={"code": "TTSSynthesisFailed", "message": str(exc)}, - ) from exc except MediaKitError as exc: - # TTS 合成失败 / 音色无权访问 → 400/403;MediaKit 提交失败 → 502 + # 音色无权访问 → 403;参数无效 → 400;MediaKit 提交失败 → 502 status_code = 502 if exc.code in ("VoiceForbidden",): status_code = 403 elif exc.code in ("InvalidInput", "TTSInvalidParam", "VoiceNotReady"): status_code = 400 - elif exc.code == "TTSSynthesisFailed": - status_code = 502 raise HTTPException( status_code=status_code, detail={ @@ -187,13 +174,13 @@ def cancel_lipsync_job( current_user: AuthenticatedUser = Depends(get_current_user), svc: LipsyncService = Depends(_get_service), ): - """取消对口型任务(仅 pending/submitted 状态可取消).""" + """取消对口型任务(仅 pending/tts_processing/submitted 状态可取消).""" job = svc.cancel_job(job_id, current_user.user.id) if job is None: raise HTTPException(status_code=404, detail="任务不存在") if job.status != "cancelled": raise HTTPException( status_code=400, - detail=f"任务状态 {job.status} 不可取消,仅 pending/submitted 可取消", + detail=f"任务状态 {job.status} 不可取消,仅 pending/tts_processing/submitted 可取消", ) return job diff --git a/apps/api/app/services/lipsync_service.py b/apps/api/app/services/lipsync_service.py index 8776ba6f5..385675980 100644 --- a/apps/api/app/services/lipsync_service.py +++ b/apps/api/app/services/lipsync_service.py @@ -25,6 +25,9 @@ from app.services.mediakit_client import ( MediaKitError, get_mediakit_client, ) + +# Celery 异步任务:TTS 合成 + MediaKit 提交(#lipsync-speed-optimization) +from app.tasks.lipsync_tts import tts_synthesize_and_submit from sqlalchemy.orm import Session from packages.adapters.sqlalchemy_impl.models import LipsyncJobModel @@ -154,41 +157,31 @@ class LipsyncService: enable_video_loop: bool = False, project_id: str = "", ) -> LipsyncJobModel: - """创建对口型任务并提交到 MediaKit. + """创建对口型任务. 两种输入模式: - - TTS 直生:voice_id + script_text(audio_url 留空),后端先合成音频 + - TTS 直生:voice_id + script_text(audio_url 留空) + → 先创建 DB 记录(状态 tts_processing),再 dispatch Celery 异步任务 + 执行 TTS 合成 + MediaKit 提交。API 响应 <1s。 - 直接音频:提供 audio_url + → 同步提交 MediaKit,状态直接设为 submitted。 Raises: - MediaKitError: TTS 合成或 MediaKit 提交失败 + MediaKitError: 参数校验失败或 MediaKit 提交失败(仅直接音频模式) """ - # 0. TTS 直生模式:先合成音频(在创建 DB 记录之前完成,失败直接抛出) + # 0. 输入校验 if not audio_url: if not (voice_id and script_text): raise MediaKitError( "必须提供 audio_url 或 voice_id+script_text", code="InvalidInput", ) - # 预合成:用临时 job_id 命名 OSS 对象 - pre_job_id = str(uuid.uuid4()) - audio_url = self._synthesize_and_persist_audio( - user_id=user_id, - job_id=pre_job_id, - voice_id=voice_id, - script_text=script_text, - speed=speed, - emotion=emotion, - ) - - # #1839: 私有桶 OSS 的裸 URL / 即将过期的短预签名会让 MediaKit GPU worker 拉取时 403, - # 提交前统一重签长有效期;外部临时 URL(CosyVoice/MediaKit)原样透传。 - video_url = self._sign_media_url(video_url) - if audio_url: - audio_url = self._sign_media_url(audio_url) + # TTS 模式:在 HTTP 请求中同步校验音色归属,快速失败 + self._resolve_voice_id(voice_id, user_id) # 1. 创建数据库记录 job_id = str(uuid.uuid4()) + is_tts_mode = not bool(audio_url) job = LipsyncJobModel( id=job_id, user_id=user_id, @@ -200,28 +193,53 @@ class LipsyncService: script_text=script_text or "", speed=speed, emotion=normalize_emotion(emotion), - status="pending", + status="tts_processing" if is_tts_mode else "pending", ) self.db.add(job) self.db.flush() - # 3. 提交到 MediaKit - try: - result = self.client.submit_lipsync( - video_url=video_url, - audio_url=audio_url, - enable_video_loop=enable_video_loop, - client_token=job_id, # 幂等控制 - ) - job.mediakit_task_id = result["task_id"] - job.status = "submitted" - job.submitted_at = datetime.now(timezone.utc) - except MediaKitError as exc: - job.status = "failed" - job.error_message = str(exc) - job.error_code = exc.code - logger.error("提交对口型任务失败: %s", exc) - raise + if is_tts_mode: + # 2a. TTS 模式:dispatch Celery 异步任务处理 TTS 合成 + MediaKit 提交 + try: + tts_synthesize_and_submit.apply_async( + args=( + job_id, + user_id, + voice_id, + script_text, + speed, + normalize_emotion(emotion), + ) + ) + except Exception: + logger.warning( + "Celery 任务提交失败,TTS 任务已创建但未触发执行: %s", + job_id, + exc_info=True, + ) + else: + # 2b. 直接音频模式:同步签名并提交 MediaKit + video_url = self._sign_media_url(video_url) + if audio_url: + audio_url = self._sign_media_url(audio_url) + job.audio_url = audio_url + + try: + result = self.client.submit_lipsync( + video_url=video_url, + audio_url=audio_url, + enable_video_loop=enable_video_loop, + client_token=job_id, + ) + job.mediakit_task_id = result["task_id"] + job.status = "submitted" + job.submitted_at = datetime.now(timezone.utc) + except MediaKitError as exc: + job.status = "failed" + job.error_message = str(exc) + job.error_code = exc.code + logger.error("提交对口型任务失败: %s", exc) + raise self.db.commit() self.db.refresh(job) @@ -360,12 +378,12 @@ class LipsyncService: # ── 取消任务 ────────────────────────────────────────────────────────── def cancel_job(self, job_id: str, user_id: str) -> Optional[LipsyncJobModel]: - """取消任务(仅 pending/submitted 状态可取消).""" + """取消任务(仅 pending/tts_processing/submitted 状态可取消).""" job = self.get_job(job_id, user_id) if job is None: return None - if job.status in ("pending", "submitted"): + if job.status in ("pending", "tts_processing", "submitted"): job.status = "cancelled" job.updated_at = datetime.now(timezone.utc) self.db.commit() diff --git a/apps/api/app/tasks/lipsync_tts.py b/apps/api/app/tasks/lipsync_tts.py new file mode 100644 index 000000000..45e1ed11e --- /dev/null +++ b/apps/api/app/tasks/lipsync_tts.py @@ -0,0 +1,166 @@ +"""AI 数字人对口型 TTS 异步任务 — 将 TTS 合成从 HTTP 请求移至 Celery 后台执行. + +优化目标:将 create_job 的 API 响应时间从 6~35s 降到 <1s。 +任务流程: +1. 创建新 DB session,加载 job 记录 +2. 调用 CosyVoice 合成音频 +3. 下载音频并转存到自家 OSS +4. 更新 job 的 audio_url +5. 签名 URL 并提交到 MediaKit +6. 更新 job 状态为 submitted +7. 异常时标记 job 为 failed +""" + +import io +import logging +from datetime import datetime, timezone + +from apps.worker.worker_app.celery_app import celery_app + +logger = logging.getLogger(__name__) + + +@celery_app.task( + bind=True, + name="lipsync_tts.synthesize_and_submit", + max_retries=2, + default_retry_delay=30, +) +def tts_synthesize_and_submit( + self, + job_id: str, + user_id: str, + voice_id: str, + script_text: str, + speed: float, + emotion: str, +): + """异步执行 TTS 合成 + OSS 转存 + MediaKit 提交. + + 在 Celery worker 中运行,不阻塞 HTTP 请求。 + """ + from app.services.mediakit_client import MediaKitError, get_mediakit_client + from sqlalchemy.orm import Session as DBSession + + from packages.adapters.sqlalchemy_impl.database import SessionLocal + from packages.adapters.sqlalchemy_impl.models import LipsyncJobModel + from packages.application.cosyvoice_service import CosyVoiceError, CosyVoiceService + from packages.shared.storage import get_shared_storage_service + from packages.shared.url_security import safe_download_bytes + + db: DBSession = SessionLocal() + try: + job = ( + db.query(LipsyncJobModel) + .filter( + LipsyncJobModel.id == job_id, + LipsyncJobModel.user_id == user_id, + ) + .first() + ) + + if job is None: + logger.error("[lipsync_tts] Job not found: job_id=%s", job_id) + return + + # 已取消的任务不再处理 + if job.status == "cancelled": + logger.info("[lipsync_tts] Job already cancelled, skipping: job_id=%s", job_id) + return + + # 1. TTS 合成 + try: + cosyvoice = CosyVoiceService() + result = cosyvoice.submit_synthesize_task( + text=script_text, + voice_id=voice_id, + speed=speed, + emotion=emotion, + ) + except CosyVoiceError as exc: + logger.error("[lipsync_tts] TTS 合成失败: job_id=%s err=%s", job_id, exc) + job.status = "failed" + job.error_message = f"TTS 合成失败: {exc}" + job.error_code = "TTSSynthesisFailed" + job.updated_at = datetime.now(timezone.utc) + db.commit() + return + except ValueError as exc: + logger.error("[lipsync_tts] TTS 参数错误: job_id=%s err=%s", job_id, exc) + job.status = "failed" + job.error_message = f"TTS 参数错误: {exc}" + job.error_code = "TTSInvalidParam" + job.updated_at = datetime.now(timezone.utc) + db.commit() + return + + temp_url = result.get("audio_url", "") + if not temp_url: + logger.error("[lipsync_tts] TTS 未返回音频 URL: job_id=%s", job_id) + job.status = "failed" + job.error_message = "TTS 未返回音频 URL" + job.error_code = "TTSNoAudio" + job.updated_at = datetime.now(timezone.utc) + db.commit() + return + + # 2. 下载并转存到自家 OSS + try: + audio_data = safe_download_bytes( + temp_url, + purpose="lipsync_tts_audio", + allowed_mime_types=("audio/mpeg", "audio/mp3", "audio/wav", "audio/mp4", "audio/x-m4a"), + timeout=60.0, + ) + storage = get_shared_storage_service() + storage_key = f"lipsync-tts/{user_id}/{job_id}.mp3" + permanent_url = storage.upload_file(io.BytesIO(audio_data), storage_key, content_type="audio/mpeg") + logger.info("[lipsync_tts] TTS 音频已转存 OSS: job_id=%s key=%s", job_id, storage_key) + job.audio_url = permanent_url + except Exception as exc: + logger.warning("[lipsync_tts] TTS 音频转存 OSS 失败,回退临时 URL: job_id=%s err=%s", job_id, exc) + job.audio_url = temp_url + + db.commit() + + # 3. 签名 URL 并提交到 MediaKit + from app.services.lipsync_service import LipsyncService + + temp_service = LipsyncService.__new__(LipsyncService) + audio_url = temp_service._sign_media_url(job.audio_url) + video_url = temp_service._sign_media_url(job.video_url) + + client = get_mediakit_client() + try: + mk_result = client.submit_lipsync( + video_url=video_url, + audio_url=audio_url, + enable_video_loop=job.enable_video_loop, + client_token=job_id, + ) + job.mediakit_task_id = mk_result["task_id"] + job.status = "submitted" + job.submitted_at = datetime.now(timezone.utc) + logger.info("[lipsync_tts] 已提交 MediaKit: job_id=%s task_id=%s", job_id, mk_result["task_id"]) + except MediaKitError as exc: + job.status = "failed" + job.error_message = str(exc) + job.error_code = exc.code + logger.error("[lipsync_tts] 提交 MediaKit 失败: job_id=%s err=%s", job_id, exc) + + db.commit() + + except Exception: + logger.exception("[lipsync_tts] 未预期的异常: job_id=%s", job_id) + try: + job = db.query(LipsyncJobModel).filter(LipsyncJobModel.id == job_id).first() + if job and job.status not in ("cancelled", "failed", "completed"): + job.status = "failed" + job.error_message = "TTS 异步任务执行异常" + job.error_code = "AsyncTaskError" + job.updated_at = datetime.now(timezone.utc) + db.commit() + except Exception: + logger.exception("[lipsync_tts] 回写失败状态时异常: job_id=%s", job_id) + finally: + db.close() diff --git a/apps/worker/worker_app/celery_app.py b/apps/worker/worker_app/celery_app.py index 34d40ad09..d10c81c74 100755 --- a/apps/worker/worker_app/celery_app.py +++ b/apps/worker/worker_app/celery_app.py @@ -37,6 +37,7 @@ celery_app.conf.imports = ( "worker_app.tasks._startup", "apps.worker.video_processing.dedup", "worker_app.tasks.cleanup", + "apps.api.app.tasks.lipsync_tts", ) # Celery Beat 定时任务调度 diff --git a/tests/unit/test_ai_avatar_emotion_tts_lipsync.py b/tests/unit/test_ai_avatar_emotion_tts_lipsync.py index 8f80e2481..74aeebc68 100644 --- a/tests/unit/test_ai_avatar_emotion_tts_lipsync.py +++ b/tests/unit/test_ai_avatar_emotion_tts_lipsync.py @@ -132,14 +132,8 @@ def _lipsync_service_with_mocks(): def test_create_job_tts_direct_mode_synthesizes_audio(): svc, client, cosy = _lipsync_service_with_mocks() - with ( - patch("app.services.lipsync_service.get_shared_storage_service") as storage_patch, - patch("app.services.lipsync_service.safe_download_bytes") as dl_patch, - ): - storage = MagicMock() - storage.upload_file.return_value = "https://oss/tts.mp3" - storage_patch.return_value = storage - dl_patch.return_value = b"FAKEAUDIO" + with patch("app.services.lipsync_service.tts_synthesize_and_submit") as mock_task: + mock_task.apply_async.return_value = MagicMock(id="celery-task-123") job = svc.create_job( user_id="user-1", @@ -150,19 +144,16 @@ def test_create_job_tts_direct_mode_synthesizes_audio(): emotion="兴奋", ) - # 调了 TTS 合成,带 speed/emotion - cosy.submit_synthesize_task.assert_called_once() - _, kwargs = cosy.submit_synthesize_task.call_args - assert kwargs["speed"] == 1.2 - assert kwargs["emotion"] == "excited" - assert kwargs["voice_id"] == "cosy-v1" - # MediaKit 用合成后的 OSS 音频 URL 提交 - _, submit_kwargs = client.submit_lipsync.call_args - assert submit_kwargs["audio_url"] == "https://oss/tts.mp3" - assert submit_kwargs["video_url"] == "https://oss/person.mp4" - # DB 记录了 TTS 字段 + # v4: TTS 模式下 create_job 返回 tts_processing 状态,dispatch Celery 任务 + assert job.status == "tts_processing" assert job.emotion == "excited" assert job.speed == 1.2 + # 不直接调用 CosyVoice(由 Celery 任务处理) + cosy.submit_synthesize_task.assert_not_called() + # 不直接提交 MediaKit(由 Celery 任务处理) + client.submit_lipsync.assert_not_called() + # dispatch 了 Celery 任务 + mock_task.apply_async.assert_called_once() def test_create_job_direct_audio_mode_skips_tts(): @@ -180,23 +171,25 @@ def test_create_job_direct_audio_mode_skips_tts(): def test_create_job_tts_failure_raises(): - from app.services.mediakit_client import MediaKitError - - from packages.application.cosyvoice_service import CosyVoiceError - + """v4: TTS 模式下 create_job 不再同步失败,而是 dispatch Celery 任务。 + TTS 合成失败由 Celery 任务内部处理并更新 job 状态。""" svc, client, cosy = _lipsync_service_with_mocks() - cosy.submit_synthesize_task.side_effect = CosyVoiceError("Arrearage") - with pytest.raises(MediaKitError) as exc: - svc.create_job( + with patch("app.services.lipsync_service.tts_synthesize_and_submit") as mock_task: + mock_task.apply_async.return_value = MagicMock(id="celery-task-456") + + job = svc.create_job( user_id="user-1", video_url="https://oss/person.mp4", voice_id="v-1", script_text="文本", ) - assert exc.value.code == "TTSSynthesisFailed" - # TTS 失败不应提交 MediaKit + + # create_job 成功返回 tts_processing,不直接调用 TTS + assert job.status == "tts_processing" + cosy.submit_synthesize_task.assert_not_called() client.submit_lipsync.assert_not_called() + mock_task.apply_async.assert_called_once() # ── refresh 同步中间状态 ──────────────────────────────────────────────── diff --git a/tests/unit/test_lipsync_routes.py b/tests/unit/test_lipsync_routes.py index 2e606a98e..a720e6dc3 100644 --- a/tests/unit/test_lipsync_routes.py +++ b/tests/unit/test_lipsync_routes.py @@ -208,24 +208,21 @@ class TestLipsyncServiceUnit: """Service 层单元测试(纯 mock,不依赖数据库)— #1809 更新.""" def test_create_job_success(self, mock_mediakit, mock_cosyvoice): - """v3: TTS 直生——service 内部 submit_synthesize_task 合成后转存 OSS,再提交 MediaKit.""" + """TTS 直生——v4 异步模式:create_job 只创建 DB 记录 + dispatch Celery 任务.""" from app.services.lipsync_service import LipsyncService mock_db = MagicMock() mock_repo = MagicMock() mock_repo.get.return_value = None # 预置音色,原样返回 voice_id - with ( - patch("app.services.lipsync_service.get_shared_storage_service") as storage_patch, - patch("app.services.lipsync_service.safe_download_bytes") as dl_patch, - ): - storage_patch.return_value.upload_file.return_value = "https://my-oss/tts.mp3" - dl_patch.return_value = b"audio-bytes" + with patch( + "app.services.lipsync_service.tts_synthesize_and_submit" + ) as mock_task: + mock_task.apply_async.return_value = MagicMock(id="celery-task-123") svc = LipsyncService( mock_db, client=mock_mediakit, - cosyvoice_service=mock_cosyvoice, voice_clone_repo=mock_repo, ) @@ -238,32 +235,23 @@ class TestLipsyncServiceUnit: emotion="兴奋", ) - assert job.status == "submitted" - assert job.mediakit_task_id == "mk-task-123" - # TTS 直生走 submit_synthesize_task,带语速/情绪 - mock_cosyvoice.submit_synthesize_task.assert_called_once() - _, kwargs = mock_cosyvoice.submit_synthesize_task.call_args - assert kwargs["text"] == "大家好,欢迎来到直播间" - assert kwargs["voice_id"] == "longxiaochun_v3" - assert kwargs["speed"] == 1.2 - assert kwargs["emotion"] == "excited" # 兴奋→excited + # TTS 模式:异步返回,状态为 tts_processing + assert job.status == "tts_processing" + assert not job.audio_url # TTS 音频尚未合成(默认空字符串) + # 不直接调用 CosyVoice + mock_cosyvoice.submit_synthesize_task.assert_not_called() + # dispatch 了 Celery 任务 + mock_task.apply_async.assert_called_once() # job 记录透传字段 assert job.speed == 1.2 assert job.emotion == "excited" - # MediaKit 用转存后的 OSS audio_url - call_kwargs = mock_mediakit.submit_lipsync.call_args - assert call_kwargs.kwargs["audio_url"] == "https://my-oss/tts.mp3" + # MediaKit 尚未提交(由 Celery 任务处理) + mock_mediakit.submit_lipsync.assert_not_called() def test_create_job_tts_failure(self, mock_mediakit): - """v3: TTS 合成失败时,CosyVoiceError 被包装为 MediaKitError(TTSSynthesisFailed), - 在建 DB 记录之前抛出,不提交 MediaKit。""" + """v4: TTS 模式下 create_job 不再同步失败,而是 dispatch Celery 任务。 + TTS 合成失败由 Celery 任务内部处理(见 test_lipsync_speed_optimization.py)。""" from app.services.lipsync_service import LipsyncService - from app.services.mediakit_client import MediaKitError - - from packages.application.cosyvoice_service import CosyVoiceError - - mock_cosyvoice = MagicMock() - mock_cosyvoice.submit_synthesize_task.side_effect = CosyVoiceError("Arrearage 欠费") mock_db = MagicMock() mock_repo = MagicMock() @@ -272,38 +260,41 @@ class TestLipsyncServiceUnit: svc = LipsyncService( mock_db, client=mock_mediakit, - cosyvoice_service=mock_cosyvoice, voice_clone_repo=mock_repo, ) - with pytest.raises(MediaKitError) as exc_info: - svc.create_job( + with patch( + "app.services.lipsync_service.tts_synthesize_and_submit" + ) as mock_task: + mock_task.apply_async.return_value = MagicMock(id="celery-task-456") + + job = svc.create_job( user_id="user-1", video_url="https://example.com/video.mp4", voice_id="longxiaochun_v3", script_text="测试文本", ) - assert exc_info.value.code == "TTSSynthesisFailed" - # 不应提交到 MediaKit + # TTS 模式下 create_job 成功返回,状态为 tts_processing + assert job.status == "tts_processing" mock_mediakit.submit_lipsync.assert_not_called() + mock_task.apply_async.assert_called_once() - def test_create_job_api_failure(self, mock_mediakit, mock_cosyvoice): - """MediaKit 提交失败.""" + def test_create_job_api_failure(self, mock_mediakit): + """MediaKit 提交失败(直传音频模式同步触发).""" from app.services.lipsync_service import LipsyncService from app.services.mediakit_client import MediaKitError mock_mediakit.submit_lipsync.side_effect = MediaKitError("API 调用失败", code="SubmitFailed") mock_db = MagicMock() - svc = LipsyncService(mock_db, client=mock_mediakit, cosyvoice_service=mock_cosyvoice) + svc = LipsyncService(mock_db, client=mock_mediakit) with pytest.raises(MediaKitError, match="API 调用失败"): svc.create_job( user_id="user-1", video_url="https://example.com/video.mp4", - voice_id="longxiaochun_v3", - script_text="测试文本", + audio_url="https://example.com/audio.mp3", ) def test_get_job_delegates_to_db(self, mock_mediakit, mock_cosyvoice): @@ -434,24 +425,21 @@ class TestLipsyncServiceUnit: assert result.status == "completed" def test_create_job_stores_tts_audio_url(self, mock_mediakit, mock_cosyvoice): - """v3: TTS 直生模式下 job.audio_url 为转存到自家 OSS 的永久地址.""" + """v4: TTS 模式下 create_job 返回 tts_processing 状态,audio_url 尚未设置(由 Celery 任务处理).""" from app.services.lipsync_service import LipsyncService mock_db = MagicMock() mock_repo = MagicMock() mock_repo.get.return_value = None - with ( - patch("app.services.lipsync_service.get_shared_storage_service") as storage_patch, - patch("app.services.lipsync_service.safe_download_bytes") as dl_patch, - ): - storage_patch.return_value.upload_file.return_value = "https://my-oss/permanent.mp3" - dl_patch.return_value = b"audio-bytes" + with patch( + "app.services.lipsync_service.tts_synthesize_and_submit" + ) as mock_task: + mock_task.apply_async.return_value = MagicMock(id="celery-task-789") svc = LipsyncService( mock_db, client=mock_mediakit, - cosyvoice_service=mock_cosyvoice, voice_clone_repo=mock_repo, ) job = svc.create_job( @@ -461,8 +449,11 @@ class TestLipsyncServiceUnit: script_text="这是一段测试文本", ) - # job.audio_url 是转存 OSS 后的永久地址 - assert job.audio_url == "https://my-oss/permanent.mp3" + # TTS 模式下 create_job 返回 tts_processing 状态 + assert job.status == "tts_processing" + # audio_url 尚未设置(由 Celery 任务异步处理),模型默认为空字符串 + assert not job.audio_url + mock_task.apply_async.assert_called_once() def test_create_job_direct_audio_skips_tts(self, mock_mediakit, mock_cosyvoice): """v3: 直接音频模式(传 audio_url)不触发 TTS,原样把 audio_url 提交 MediaKit.""" @@ -548,9 +539,9 @@ class TestErrorHandling: assert exc_info.value.code == "VoiceNotReady" def test_tts_value_error_mapped_to_invalid_param(self, mock_mediakit): - """v3: CosyVoice 抛 ValueError(参数无效)被包装为 TTSInvalidParam(路由映射 400).""" + """v4: TTS 模式下 create_job 不再同步调用 CosyVoice, + 而是 dispatch Celery 任务。ValueError 由 Celery 任务内部处理。""" from app.services.lipsync_service import LipsyncService - from app.services.mediakit_client import MediaKitError mock_cosyvoice = MagicMock() mock_cosyvoice.submit_synthesize_task.side_effect = ValueError("voice_id 为空") @@ -562,18 +553,27 @@ class TestErrorHandling: svc = LipsyncService( mock_db, client=mock_mediakit, - cosyvoice_service=mock_cosyvoice, voice_clone_repo=mock_repo, ) - with pytest.raises(MediaKitError) as exc_info: - svc.create_job( + with patch( + "app.services.lipsync_service.tts_synthesize_and_submit" + ) as mock_task: + mock_task.apply_async.return_value = MagicMock(id="celery-task-789") + + # TTS 模式下 create_job 不再同步失败 + job = svc.create_job( user_id="user-1", video_url="https://example.com/video.mp4", voice_id="some-voice", script_text="test", ) - assert exc_info.value.code == "TTSInvalidParam" + + # 确认返回 tts_processing 状态 + assert job.status == "tts_processing" + # TTS 合成由 Celery 任务处理,不直接调用 CosyVoice + mock_cosyvoice.submit_synthesize_task.assert_not_called() + mock_task.apply_async.assert_called_once() def test_missing_both_inputs_raises_invalid_input(self, mock_mediakit, mock_cosyvoice): """v3: 既无 audio_url 又无 voice_id+script_text 时抛 InvalidInput(路由映射 400).""" diff --git a/tests/unit/test_lipsync_speed_optimization.py b/tests/unit/test_lipsync_speed_optimization.py new file mode 100644 index 000000000..67752d1c7 --- /dev/null +++ b/tests/unit/test_lipsync_speed_optimization.py @@ -0,0 +1,623 @@ +"""AI 数字人口型视频生成速度优化 — 单元测试. + +验证两个优化点: +1. FFmpeg 编码 preset 从 fast 改为 veryfast(提速 30~50%) +2. TTS 合成从同步改为 Celery 异步任务(API 响应从 6~35s 降到 <1s) + +Issue: lipsync-speed-optimization +""" + +import os +from unittest.mock import MagicMock, patch + +import pytest + +os.environ.setdefault("JWT_SECRET_KEY", "dev-secret-key-for-testing") + + +# ═══════════════════════════════════════════════════════════════════════════════ +# 优化1: FFmpeg 编码提速 — preset veryfast +# ═══════════════════════════════════════════════════════════════════════════════ + + +class TestFFmpegPresetOptimization: + """验证 FFmpeg 编码命令从 -preset fast 改为 -preset veryfast.""" + + def test_preset_is_veryfast(self): + """_build_ffmpeg_command 输出必须包含 -preset veryfast.""" + from app.services.ai_avatar_render_service import AiAvatarRenderService + + svc = AiAvatarRenderService.__new__(AiAvatarRenderService) + cmd = svc._build_ffmpeg_command( + input_video="https://example.com/video.mp4", + b_roll_segments=[], + filter_complex="", + final_label=None, + output_path="/tmp/output.mp4", + ) + assert "-preset veryfast" in cmd, f"期望 -preset veryfast,实际命令: {cmd}" + + def test_preset_veryfast_with_filter(self): + """带滤镜场景下也必须使用 veryfast.""" + from app.services.ai_avatar_render_service import AiAvatarRenderService + + svc = AiAvatarRenderService.__new__(AiAvatarRenderService) + cmd = svc._build_ffmpeg_command( + input_video="https://example.com/video.mp4", + b_roll_segments=[], + filter_complex="overlay=0:0", + final_label="[v]", + output_path="/tmp/output.mp4", + ) + assert "-preset veryfast" in cmd + assert "-filter_complex" in cmd + + def test_preset_not_fast(self): + """确保不再使用旧的 -preset fast.""" + from app.services.ai_avatar_render_service import AiAvatarRenderService + + svc = AiAvatarRenderService.__new__(AiAvatarRenderService) + cmd = svc._build_ffmpeg_command( + input_video="https://example.com/video.mp4", + b_roll_segments=[], + filter_complex="", + final_label=None, + output_path="/tmp/output.mp4", + ) + # 确保是 veryfast 而不是 fast + assert "-preset veryfast" in cmd + # 排除 "fast" 单独出现(veryfast 包含 fast 子串,需精确判断) + parts = cmd.split() + preset_idx = parts.index("-preset") + assert parts[preset_idx + 1] == "veryfast" + + +# ═══════════════════════════════════════════════════════════════════════════════ +# 优化2: TTS 合成 Celery 异步化 +# ═══════════════════════════════════════════════════════════════════════════════ + + +def _make_service_with_mocks(): + """构造 LipsyncService 测试实例及 mock 依赖.""" + from app.services.lipsync_service import LipsyncService + + db = MagicMock() + client = MagicMock() + client.is_available = True + client.submit_lipsync.return_value = { + "success": True, + "task_id": "mk-1", + "request_id": "req-1", + } + cosy = MagicMock() + cosy.submit_synthesize_task.return_value = { + "audio_url": "https://tts/raw.mp3", + "request_id": "tts-req", + "audio_duration": 3.0, + } + svc = LipsyncService(db, client=client, cosyvoice_service=cosy, voice_clone_repo=MagicMock()) + # _resolve_voice_id 默认原样返回(repo.get 返回 None) + svc._voice_clone_repo.get.return_value = None + return svc, client, cosy + + +class TestCreateJobAsyncTTS: + """验证 TTS 模式改为 Celery 异步后的行为.""" + + def test_tts_mode_returns_tts_processing_status(self): + """TTS 模式下 create_job 立即返回,状态为 tts_processing.""" + svc, client, cosy = _make_service_with_mocks() + + with patch("app.services.lipsync_service.tts_synthesize_and_submit") as mock_task: + mock_task.apply_async = MagicMock() + + job = svc.create_job( + user_id="user-1", + video_url="https://example.com/video.mp4", + voice_id="longxiaochun_v3", + script_text="大家好", + speed=1.0, + emotion="", + ) + + assert job.status == "tts_processing" + + def test_tts_mode_dispatches_celery_task(self): + """TTS 模式必须 dispatch Celery 异步任务.""" + svc, client, cosy = _make_service_with_mocks() + + with patch("app.services.lipsync_service.tts_synthesize_and_submit") as mock_task: + mock_task.apply_async = MagicMock() + + svc.create_job( + user_id="user-1", + video_url="https://example.com/video.mp4", + voice_id="v-1", + script_text="测试文本", + ) + + mock_task.apply_async.assert_called_once() + call_kwargs = mock_task.apply_async.call_args + args = call_kwargs.kwargs.get("args") or call_kwargs[1].get("args", call_kwargs[0][0] if call_kwargs[0] else ()) + assert args[1] == "user-1" # user_id + assert args[2] == "v-1" # voice_id + assert args[3] == "测试文本" # script_text + + def test_tts_mode_celery_dispatch_failure_still_creates_job(self): + """Celery dispatch 失败时,job 记录已创建,状态保持 tts_processing.""" + svc, client, cosy = _make_service_with_mocks() + + with patch("app.services.lipsync_service.tts_synthesize_and_submit") as mock_task: + mock_task.apply_async = MagicMock(side_effect=Exception("Celery broker down")) + + job = svc.create_job( + user_id="user-1", + video_url="https://example.com/video.mp4", + voice_id="v-1", + script_text="测试文本", + ) + + # job 已创建 + assert job is not None + assert job.status == "tts_processing" + # MediaKit 未被调用 + client.submit_lipsync.assert_not_called() + + def test_tts_mode_voice_validation_still_sync(self): + """TTS 模式下音色校验仍在 HTTP 请求中同步执行.""" + from app.services.mediakit_client import MediaKitError + + svc, client, cosy = _make_service_with_mocks() + # 模拟音色属于其他用户 + other_profile = MagicMock() + other_profile.user_id = "user-other" + svc._voice_clone_repo.get.return_value = other_profile + + with patch("app.tasks.lipsync_tts.tts_synthesize_and_submit"): + with pytest.raises(MediaKitError) as exc: + svc.create_job( + user_id="user-1", + video_url="https://example.com/video.mp4", + voice_id="clone-profile-id", + script_text="测试", + ) + assert exc.value.code == "VoiceForbidden" + + def test_tts_mode_missing_input_raises_immediately(self): + """缺少 voice_id 或 script_text 时立即报错,不 dispatch Celery 任务.""" + from app.services.mediakit_client import MediaKitError + + svc, client, cosy = _make_service_with_mocks() + + with patch("app.tasks.lipsync_tts.tts_synthesize_and_submit") as mock_task: + mock_task.delay = MagicMock() + + with pytest.raises(MediaKitError) as exc: + svc.create_job( + user_id="user-1", + video_url="https://example.com/video.mp4", + # 缺少 voice_id 和 script_text + ) + assert exc.value.code == "InvalidInput" + + # Celery 任务未被 dispatch + mock_task.delay.assert_not_called() + # TTS 和 MediaKit 均未调用 + cosy.submit_synthesize_task.assert_not_called() + client.submit_lipsync.assert_not_called() + + +class TestCreateJobDirectAudio: + """验证直接音频模式不受异步化影响.""" + + def test_direct_audio_still_submits_synchronously(self): + """直接音频模式仍然同步提交 MediaKit,状态为 submitted.""" + svc, client, cosy = _make_service_with_mocks() + + with patch("app.tasks.lipsync_tts.tts_synthesize_and_submit") as mock_task: + mock_task.delay = MagicMock() + + job = svc.create_job( + user_id="user-1", + video_url="https://example.com/video.mp4", + audio_url="https://example.com/audio.mp3", + ) + + assert job.status == "submitted" + assert job.mediakit_task_id == "mk-1" + client.submit_lipsync.assert_called_once() + # TTS Celery 任务不应被调用 + mock_task.delay.assert_not_called() + + def test_direct_audio_skips_tts(self): + """直接音频模式不调用 CosyVoice TTS.""" + svc, client, cosy = _make_service_with_mocks() + + job = svc.create_job( + user_id="user-1", + video_url="https://example.com/video.mp4", + audio_url="https://example.com/audio.mp3", + ) + + cosy.submit_synthesize_task.assert_not_called() + call_kwargs = client.submit_lipsync.call_args + assert call_kwargs.kwargs["audio_url"] == "https://example.com/audio.mp3" + + +class TestCancelJobTtsProcessing: + """验证 tts_processing 状态的任务可以被取消.""" + + def test_cancel_tts_processing(self): + """tts_processing 状态的任务可以成功取消.""" + svc, client, cosy = _make_service_with_mocks() + + mock_job = MagicMock() + mock_job.status = "tts_processing" + mock_job.id = "job-1" + svc.get_job = MagicMock(return_value=mock_job) + + result = svc.cancel_job("job-1", "user-1") + assert result.status == "cancelled" + + def test_cancel_pending_still_works(self): + """pending 状态仍可取消.""" + svc, client, cosy = _make_service_with_mocks() + + mock_job = MagicMock() + mock_job.status = "pending" + mock_job.id = "job-1" + svc.get_job = MagicMock(return_value=mock_job) + + result = svc.cancel_job("job-1", "user-1") + assert result.status == "cancelled" + + def test_cancel_submitted_still_works(self): + """submitted 状态仍可取消.""" + svc, client, cosy = _make_service_with_mocks() + + mock_job = MagicMock() + mock_job.status = "submitted" + mock_job.id = "job-1" + svc.get_job = MagicMock(return_value=mock_job) + + result = svc.cancel_job("job-1", "user-1") + assert result.status == "cancelled" + + +# ═══════════════════════════════════════════════════════════════════════════════ +# lipsync_tts.py — Celery 异步任务单元测试 +# ═══════════════════════════════════════════════════════════════════════════════ + +import sys +import types + + +class _FakeQuery: + """模拟 SQLAlchemy query.filter().first() 链式调用.""" + + def __init__(self, job): + self._job = job + + def filter(self, *args, **kwargs): + return self + + def first(self): + return self._job + + +def _make_fake_job(**kwargs): + """构造可 setattr 的 job 记录.""" + job = MagicMock() + job.id = kwargs.get("job_id", "job-1") + job.user_id = kwargs.get("user_id", "user-1") + job.status = kwargs.get("status", "tts_processing") + job.audio_url = kwargs.get("audio_url", "") + job.video_url = kwargs.get("video_url", "https://oss/video.mp4") + job.mediakit_task_id = kwargs.get("mediakit_task_id", "") + job.enable_video_loop = kwargs.get("enable_video_loop", False) + job.error_code = "" + job.error_message = "" + job.submitted_at = None + job.updated_at = None + return job + + +def _build_session(job): + """构造 mock DB session + factory. 返回 (session, factory_patch_ctx_value).""" + session = MagicMock() + session.query.return_value = _FakeQuery(job) + session.commit = MagicMock() + session.close = MagicMock() + factory = MagicMock(return_value=session) + return session, factory + + +def _apply_all_patches( + *, + job=None, + cosyvoice_service=None, + cosyvoice_side_effect=None, + cosyvoice_error=None, + download_bytes=b"AUDIO", + download_error=None, + storage=None, + mk_client=None, + mk_submit_return=None, + mk_submit_error=None, +): + """统一构造测试需要的 patch 列表. + + lipsync_tts.run() 在函数体内部懒 import 多个模块,通过 sys.modules 注入 + 伪造包路径避免真实导入;对存在的模块用 patch() 替换返回值/side_effect。 + """ + # 构造不存在的 database 模块 + fake_db_mod = types.ModuleType("packages.adapters.sqlalchemy_impl.database") + session, factory = _build_session(job) + fake_db_mod.SessionLocal = factory + + patches = [ + patch.dict(sys.modules, {"packages.adapters.sqlalchemy_impl.database": fake_db_mod}), + patch( + "app.services.lipsync_service.LipsyncService._sign_media_url", + side_effect=lambda url: url + "?signed" if url else url, + ), + ] + + # CosyVoice + if cosyvoice_service is not None: + cosy_instance = cosyvoice_service + else: + cosy_instance = MagicMock() + if cosyvoice_side_effect is not None: + cosy_instance.submit_synthesize_task.side_effect = cosyvoice_side_effect + elif cosyvoice_error is not None: + cosy_instance.submit_synthesize_task.side_effect = cosyvoice_error + else: + cosy_instance.submit_synthesize_task.return_value = {"audio_url": "https://tts/raw.mp3"} + patches.append(patch("packages.application.cosyvoice_service.CosyVoiceService", return_value=cosy_instance)) + + # safe_download_bytes + if download_error is not None: + patches.append(patch("packages.shared.url_security.safe_download_bytes", side_effect=download_error)) + else: + patches.append(patch("packages.shared.url_security.safe_download_bytes", return_value=download_bytes)) + + # Storage + if storage is None: + storage = MagicMock() + storage.public_url = "https://oss.example.com" + storage.upload_file.return_value = "https://oss.example.com/tts.mp3" + patches.append(patch("packages.shared.storage.get_shared_storage_service", return_value=storage)) + + # MediaKit client + if mk_client is not None: + patches.append(patch("app.services.mediakit_client.get_mediakit_client", return_value=mk_client)) + else: + client = MagicMock() + if mk_submit_error is not None: + client.submit_lipsync.side_effect = mk_submit_error + else: + client.submit_lipsync.return_value = mk_submit_return or {"task_id": "mk-1"} + patches.append(patch("app.services.mediakit_client.get_mediakit_client", return_value=client)) + + return session, patches + + +class TestTtsSynthesizeAndSubmit: + """测试 Celery 任务 tts_synthesize_and_submit.run 的所有分支.""" + + def test_job_not_found_returns_early(self): + """Job 不存在 → 日志报错直接返回,不抛异常.""" + from app.tasks.lipsync_tts import tts_synthesize_and_submit + + session, patches = _apply_all_patches(job=None) + entered = [p.__enter__() for p in patches] + try: + tts_synthesize_and_submit.run("missing-job", "user-1", "v1", "你好", 1.0, "") + finally: + for p in reversed(patches): + p.__exit__(None, None, None) + + session.commit.assert_not_called() + session.close.assert_called_once() + + def test_cancelled_job_skipped(self): + """Job 已 cancelled → 跳过不处理,不调用 TTS/MediaKit.""" + from app.tasks.lipsync_tts import tts_synthesize_and_submit + + job = _make_fake_job(status="cancelled") + session, patches = _apply_all_patches(job=job) + entered = [p.__enter__() for p in patches] + try: + tts_synthesize_and_submit.run("job-1", "user-1", "v1", "你好", 1.0, "") + finally: + for p in reversed(patches): + p.__exit__(None, None, None) + + session.commit.assert_not_called() + session.close.assert_called_once() + assert job.status == "cancelled" + + def test_happy_path_tts_to_mediakit(self): + """正常流程:TTS 合成 → 下载 → OSS → 签名 → 提交 MediaKit → submitted.""" + from app.tasks.lipsync_tts import tts_synthesize_and_submit + + job = _make_fake_job(status="tts_processing") + storage = MagicMock() + storage.public_url = "https://oss.example.com" + storage.upload_file.return_value = "https://oss.example.com/lipsync-tts/u/j.mp3" + + session, patches = _apply_all_patches( + job=job, + storage=storage, + mk_submit_return={"task_id": "mk-999"}, + ) + for p in patches: + p.__enter__() + try: + tts_synthesize_and_submit.run("job-1", "user-1", "v1", "你好", 1.0, "") + finally: + for p in reversed(patches): + p.__exit__(None, None, None) + + assert job.audio_url == "https://oss.example.com/lipsync-tts/u/j.mp3" + assert job.mediakit_task_id == "mk-999" + assert job.status == "submitted" + assert job.submitted_at is not None + session.close.assert_called_once() + + def test_tts_cosyvoice_error_marks_failed(self): + """CosyVoiceError → 标记 failed,error_code=TTSSynthesisFailed.""" + from app.tasks.lipsync_tts import tts_synthesize_and_submit + + from packages.application.cosyvoice_service import CosyVoiceError + + job = _make_fake_job(status="tts_processing") + session, patches = _apply_all_patches( + job=job, + cosyvoice_error=CosyVoiceError("TTS 服务异常"), + ) + for p in patches: + p.__enter__() + try: + tts_synthesize_and_submit.run("job-1", "user-1", "v1", "你好", 1.0, "") + finally: + for p in reversed(patches): + p.__exit__(None, None, None) + + assert job.status == "failed" + assert job.error_code == "TTSSynthesisFailed" + assert "TTS 合成失败" in job.error_message + session.close.assert_called_once() + + def test_tts_value_error_marks_failed(self): + """ValueError → 标记 failed,error_code=TTSInvalidParam.""" + from app.tasks.lipsync_tts import tts_synthesize_and_submit + + job = _make_fake_job(status="tts_processing") + session, patches = _apply_all_patches( + job=job, + cosyvoice_error=ValueError("speed 参数非法"), + ) + for p in patches: + p.__enter__() + try: + tts_synthesize_and_submit.run("job-1", "user-1", "v1", "你好", 1.0, "") + finally: + for p in reversed(patches): + p.__exit__(None, None, None) + + assert job.status == "failed" + assert job.error_code == "TTSInvalidParam" + assert "TTS 参数错误" in job.error_message + session.close.assert_called_once() + + def test_tts_no_audio_url_marks_failed(self): + """TTS 返回空 audio_url → 标记 failed,error_code=TTSNoAudio.""" + from app.tasks.lipsync_tts import tts_synthesize_and_submit + + job = _make_fake_job(status="tts_processing") + cosy = MagicMock() + cosy.submit_synthesize_task.return_value = {"audio_url": ""} + session, patches = _apply_all_patches(job=job, cosyvoice_service=cosy) + for p in patches: + p.__enter__() + try: + tts_synthesize_and_submit.run("job-1", "user-1", "v1", "你好", 1.0, "") + finally: + for p in reversed(patches): + p.__exit__(None, None, None) + + assert job.status == "failed" + assert job.error_code == "TTSNoAudio" + session.close.assert_called_once() + + def test_oss_upload_failure_falls_back_to_temp_url(self): + """OSS 上传失败 → 回退临时 URL,继续提交 MediaKit.""" + from app.tasks.lipsync_tts import tts_synthesize_and_submit + + job = _make_fake_job(status="tts_processing") + storage = MagicMock() + storage.public_url = "https://oss.example.com" + storage.upload_file.side_effect = Exception("OSS 上传超时") + session, patches = _apply_all_patches( + job=job, + storage=storage, + mk_submit_return={"task_id": "mk-77"}, + ) + for p in patches: + p.__enter__() + try: + tts_synthesize_and_submit.run("job-1", "user-1", "v1", "你好", 1.0, "") + finally: + for p in reversed(patches): + p.__exit__(None, None, None) + + # 回退到临时 URL + assert job.audio_url == "https://tts/raw.mp3" + assert job.status == "submitted" + assert job.mediakit_task_id == "mk-77" + session.close.assert_called_once() + + def test_mediakit_submit_failure_marks_failed(self): + """MediaKit 提交失败(MediaKitError)→ 标记 failed.""" + from app.services.mediakit_client import MediaKitError + from app.tasks.lipsync_tts import tts_synthesize_and_submit + + job = _make_fake_job(status="tts_processing") + err = MediaKitError("GPU 不可用", code="MediaKitUnavailable") + session, patches = _apply_all_patches( + job=job, + mk_submit_error=err, + ) + for p in patches: + p.__enter__() + try: + tts_synthesize_and_submit.run("job-1", "user-1", "v1", "你好", 1.0, "") + finally: + for p in reversed(patches): + p.__exit__(None, None, None) + + assert job.status == "failed" + assert job.error_code == "MediaKitUnavailable" + session.close.assert_called_once() + + def test_top_level_exception_marks_async_task_error(self): + """顶层意外异常 → except 分支回写 failed,error_code=AsyncTaskError.""" + from app.tasks.lipsync_tts import tts_synthesize_and_submit + + job = _make_fake_job(status="tts_processing") + # 不调用 _apply_all_patches,手动构造所有 patch,让 CosyVoiceService 抛异常 + fake_db_mod = types.ModuleType("packages.adapters.sqlalchemy_impl.database") + session_mock = MagicMock() + session_mock.query.return_value = _FakeQuery(job) + session_mock.commit = MagicMock() + session_mock.close = MagicMock() + fake_db_mod.SessionLocal = MagicMock(return_value=session_mock) + + all_patches = [ + patch.dict(sys.modules, {"packages.adapters.sqlalchemy_impl.database": fake_db_mod}), + patch( + "app.services.lipsync_service.LipsyncService._sign_media_url", + side_effect=lambda url: url + "?signed" if url else url, + ), + patch( + "packages.application.cosyvoice_service.CosyVoiceService", + side_effect=RuntimeError("unexpected init failure"), + ), + patch("packages.shared.url_security.safe_download_bytes", return_value=b"AUDIO"), + patch("packages.shared.storage.get_shared_storage_service", return_value=MagicMock()), + patch("app.services.mediakit_client.get_mediakit_client", return_value=MagicMock()), + ] + for p in all_patches: + p.__enter__() + try: + tts_synthesize_and_submit.run("job-1", "user-1", "v1", "你好", 1.0, "") + finally: + for p in reversed(all_patches): + p.__exit__(None, None, None) + + assert job.status == "failed" + assert job.error_code == "AsyncTaskError" + assert "TTS 异步任务执行异常" in job.error_message + session_mock.close.assert_called_once()