"""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 import uuid 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()