diff --git a/apps/api/app/tasks/lipsync_tts.py b/apps/api/app/tasks/lipsync_tts.py index 6c9fa417e..85a478828 100644 --- a/apps/api/app/tasks/lipsync_tts.py +++ b/apps/api/app/tasks/lipsync_tts.py @@ -9,18 +9,52 @@ 5. 签名 URL 并提交到 MediaKit 6. 更新 job 状态为 submitted 7. 异常时标记 job 为 failed + +注意:使用 @shared_task 而非绑定到某个 celery_app 实例, +确保任务能被 Worker 侧 celery_app 正确注册,同时 API 侧 send_task/apply_async 仍可正常调用。 """ import io import logging from datetime import datetime, timezone +from urllib.parse import urlparse -from app.core.celery_app import celery_app +from celery import shared_task logger = logging.getLogger(__name__) +# MediaKit 预签名 URL 有效期(7天,秒),与 LipsyncService._sign_media_url 保持一致 +_MEDIAKIT_URL_TTL_SECONDS = 7 * 24 * 3600 -@celery_app.task( + +def _sign_media_url(url: str) -> str: + """对自家 OSS 私有桶 URL 重签长有效期预签名. + + - 自家 OSS URL → 重签 7 天有效期 + - 外部临时 URL → 原样透传 + - 任何异常降级原样返回,不阻断主流程 + """ + if not url: + return url + try: + from packages.shared.storage import get_shared_storage_service + + storage = get_shared_storage_service() + public_base = getattr(storage, "public_url", "") + if not isinstance(public_base, str) or not public_base: + return url + own_host = urlparse(public_base).netloc.lower() + host = urlparse(url).netloc.lower() + if not own_host or host != own_host: + return url + signed = storage.get_download_url(url, expires_seconds=_MEDIAKIT_URL_TTL_SECONDS) + return signed or url + except Exception as exc: # noqa: BLE001 + logger.warning("[lipsync_tts] URL 重签失败,原样返回: url_prefix=%s err=%s", url[:80], exc) + return url + + +@shared_task( bind=True, name="lipsync_tts.synthesize_and_submit", max_retries=2, @@ -45,7 +79,6 @@ def tts_synthesize_and_submit( 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() @@ -109,26 +142,35 @@ def tts_synthesize_and_submit( 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"), + allowed_mime_types=( + "audio/mpeg", + "audio/mp3", + "audio/wav", + "audio/mp4", + "audio/x-m4a", + ), timeout=60.0, ) + from packages.shared.storage import get_shared_storage_service + 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) + 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) + # 3. 签名 URL 并提交到 MediaKit(复用模块内 _sign_media_url,避免对 LipsyncService 的耦合) + audio_url = _sign_media_url(job.audio_url) + video_url = _sign_media_url(job.video_url) client = get_mediakit_client() try: @@ -141,7 +183,11 @@ def tts_synthesize_and_submit( 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"]) + 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) diff --git a/apps/web/src/pages/ai-avatar/api/aiAvatar.ts b/apps/web/src/pages/ai-avatar/api/aiAvatar.ts index d4d8e379b..13b9acf1b 100644 --- a/apps/web/src/pages/ai-avatar/api/aiAvatar.ts +++ b/apps/web/src/pages/ai-avatar/api/aiAvatar.ts @@ -54,7 +54,7 @@ export const createLipsyncJob = async (data: { } export const getLipsyncJob = async (id: string): Promise => { - const response = await apiClient.get(`/lipsync/jobs/${id}`) + const response = await apiClient.get(`/lipsync/jobs/${id}`, { timeout: 60000 }) return response.data } diff --git a/tests/unit/test_lipsync_speed_optimization.py b/tests/unit/test_lipsync_speed_optimization.py index 67752d1c7..9adfd44d4 100644 --- a/tests/unit/test_lipsync_speed_optimization.py +++ b/tests/unit/test_lipsync_speed_optimization.py @@ -358,7 +358,7 @@ def _apply_all_patches( patches = [ patch.dict(sys.modules, {"packages.adapters.sqlalchemy_impl.database": fake_db_mod}), patch( - "app.services.lipsync_service.LipsyncService._sign_media_url", + "app.tasks.lipsync_tts._sign_media_url", side_effect=lambda url: url + "?signed" if url else url, ), ] @@ -598,7 +598,7 @@ class TestTtsSynthesizeAndSubmit: all_patches = [ patch.dict(sys.modules, {"packages.adapters.sqlalchemy_impl.database": fake_db_mod}), patch( - "app.services.lipsync_service.LipsyncService._sign_media_url", + "app.tasks.lipsync_tts._sign_media_url", side_effect=lambda url: url + "?signed" if url else url, ), patch(