fix(lipsync): TTS异步任务改用@shared_task,解决Worker侧celery_app不注册任务导致卡tts_processing
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 2s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 44s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 44s
CI/CD Pipeline / Frontend Unit Tests (pull_request) Successful in 1m10s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 1m17s
CI/CD Pipeline / PR Build Web Image (pull_request) Successful in 1m16s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 1m25s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 1m55s
CI/CD Pipeline / Frontend Lint (pull_request) Successful in 1m58s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 2m3s
CI/CD Pipeline / Validate - Style (pull_request) Successful in 2m27s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 2m53s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 6m8s
CI/CD Pipeline / CI Gate (pull_request) Failing after 3s
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
AI Code Review / AI Code Review (pull_request) Successful in 6m17s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 3m26s
CI/CD Pipeline / Deploy Production (pull_request) Failing after 134h58m32s
CI/CD Pipeline / Build Production API Image (pull_request) Failing after 134h58m35s
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Failing after 135h4m39s
CI/CD Pipeline / Build Staging API Image (pull_request) Failing after 135h4m45s
CI/CD Pipeline / Staging E2E Tests (pull_request) Failing after 135h4m14s
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Failing after 135h4m25s
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Failing after 135h4m25s
CI/CD Pipeline / Build Staging Worker Image (pull_request) Failing after 135h4m32s
CI/CD Pipeline / Build Staging Web Image (pull_request) Failing after 135h4m32s
CI/CD Pipeline / Check push changed paths (pull_request) Failing after 135h4m34s
CI/CD Pipeline / ACR Image Cleanup (pull_request) Failing after 135h4m13s
CI/CD Pipeline / Canary Release to Production (pull_request) Failing after 134h58m19s
CI/CD Pipeline / Build Production Worker Image (pull_request) Failing after 134h58m21s
CI/CD Pipeline / Build Production Web Image (pull_request) Failing after 134h58m21s
CI/CD Pipeline / Staging API Integration Tests (pull_request) Failing after 135h4m14s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Failing after 135h39m24s

P0: lipsync_tts.py 原从 app.core.celery_app 导入 celery_app 并 @celery_app.task 装饰,
task 被注册到 API 侧 celery_app 实例上,但 Worker 进程使用的是 Worker 侧 celery_app,
导致 Worker 消费消息后找不到处理函数,对口型任务永远卡在 tts_processing 状态。

修复:
1. 改为 @shared_task(name="lipsync_tts.synthesize_and_submit") 装饰器,
   任务自动注册到当前活跃的 celery_app(Worker 启动时会正确注册),
   API 侧 apply_async 调用方式不变。
2. _sign_media_url 保持为模块内独立函数(不依赖 LipsyncService 反射),
   删除函数体开头冗余的 get_shared_storage_service 重复 import,
   仅在 OSS 上传段按需懒导入,保持 task 注册阶段无副作用。
3. 前端 aiAvatar.ts 的 getLipsyncJob 接口加 { timeout: 60000 },
   避免轮询时 10s 默认超时导致前端状态刷新失败。

测试:56 passed,lipsync_tts.py 覆盖率 79%。
This commit is contained in:
xiaoxia
2026-09-11 00:17:26 +08:00
parent 582f73c2f2
commit e85626dc69
3 changed files with 61 additions and 15 deletions
+58 -12
View File
@@ -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)
+1 -1
View File
@@ -54,7 +54,7 @@ export const createLipsyncJob = async (data: {
}
export const getLipsyncJob = async (id: string): Promise<LipsyncJob> => {
const response = await apiClient.get<LipsyncJob>(`/lipsync/jobs/${id}`)
const response = await apiClient.get<LipsyncJob>(`/lipsync/jobs/${id}`, { timeout: 60000 })
return response.data
}
@@ -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(