Files
xiaoxia-saas/apps/api/app/services/lipsync_service.py
T
xiaoxia 2fa6de29bc
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 1s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
CI/CD Pipeline / PR Build Web Image (pull_request) Successful in 30s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 33s
CI/CD Pipeline / Frontend Lint (pull_request) Successful in 1m39s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 1m39s
CI/CD Pipeline / Frontend Unit Tests (pull_request) Successful in 1m47s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 2m44s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m55s
AI Code Review / AI Code Review (pull_request) Successful in 6m31s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 8m15s
CI/CD Pipeline / Unit Tests (pull_request) Successful in 8m17s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 10m28s
CI/CD Pipeline / Validate - Style (pull_request) Has been cancelled
CI/CD Pipeline / Validate - Security (pull_request) Has been cancelled
CI/CD Pipeline / Build Production API Image (pull_request) Has been cancelled
CI/CD Pipeline / Build Production Web Image (pull_request) Has been cancelled
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been cancelled
CI/CD Pipeline / Deploy Production (pull_request) Has been cancelled
CI/CD Pipeline / Production Browser E2E (pull_request) Has been cancelled
CI/CD Pipeline / Canary Release to Production (pull_request) Has been cancelled
CI/CD Pipeline / CI Gate (pull_request) Has been cancelled
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been cancelled
CI/CD Pipeline / Staging API Integration Tests (pull_request) Failing after 89h23m49s
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Failing after 89h23m53s
CI/CD Pipeline / Build Staging Web Image (pull_request) Failing after 89h23m55s
CI/CD Pipeline / Build Staging API Image (pull_request) Failing after 89h23m33s
CI/CD Pipeline / ACR Image Cleanup (pull_request) Failing after 89h23m27s
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Failing after 89h23m31s
CI/CD Pipeline / Build Staging Worker Image (pull_request) Failing after 89h23m33s
CI/CD Pipeline / Staging E2E Tests (pull_request) Failing after 89h23m27s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Failing after 89h23m28s
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Failing after 89h23m31s
CI/CD Pipeline / Check push changed paths (pull_request) Failing after 89h24m9s
fix(lipsync): 修复Celery事务竞态导致job永远卡在tts_processing(P0)
根因:create_job() 先 apply_async() 发送 Celery 任务,再 db.commit() 提交事务。
worker 是独立进程+独立DB连接,任务<4ms就被消费,但此时 API 事务还未 commit,
worker 查询 job 返回 None → 静默 return 不重试 → job 永远卡在 tts_processing。

修复:
1. API侧(lipsync_service.py):先 db.commit()+db.refresh(job),再 apply_async() 发任务;
   投递失败/MediaKit提交失败分支也各自 commit,确保状态及时落库。
2. Worker侧(lipsync_tts.py):job not found 时改用 self.retry() 递增重试3次
   (1s/3s/7s退避),作为竞态场景的第二道防线;
   增加 autoretry_for=(OSError, ConnectionError) 自动重试网络抖动;
   max_retries 从2调整到5。

子agent在staging实锤:受影响job共5个,手工重投递后全部3秒内完成TTS+提交MediaKit,
证实TTS本身只需要2秒,卡顿完全是因为竞态。

单测:新增 TestCreateJobCommitOrder 验证 commit 在 apply_async 之前。
2026-09-12 21:50:02 +08:00

421 lines
17 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""对口型 Service — #1796 MediaKit 对口型业务逻辑, #1809 参数调整.
职责:
- 创建/查询对口型任务
- 双输入模式:TTS 直生(voice_id + script_text,内部先合成音频转存 OSS)或直接音频(audio_url)
- 调用 MediaKit 客户端提交异步任务
- 轮询更新任务状态(中间状态同步 DB,成片转存自家 OSS)
- 用户隔离(每个用户只能操作自己的任务)
"""
from __future__ import annotations
import io
import logging
import uuid
from datetime import datetime, timezone
from typing import Optional
from urllib.parse import urlparse
from app.services.mediakit_client import (
STATUS_COMPLETED,
STATUS_FAILED,
STATUS_RUNNING,
MediaKitClient,
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
from packages.application.cosyvoice_service import CosyVoiceError, normalize_emotion
from packages.shared.storage import get_shared_storage_service
from packages.shared.url_security import ALLOWED_AUDIO_MIME_TYPES, safe_download_bytes
logger = logging.getLogger(__name__)
# 传给 MediaKit GPU worker / 回给前端播放的 OSS 预签名有效期:7 天。
# MediaKit 排队 + 拉取可能延迟,私有桶裸 URL 或 1 小时短预签名都会 403,故统一重签长有效期。
MEDIAKIT_URL_TTL_SECONDS = 7 * 24 * 3600
class LipsyncService:
"""对口型任务 Service."""
def __init__(
self,
db: Session,
client: Optional[MediaKitClient] = None,
cosyvoice_service=None,
voice_clone_repo=None,
):
self.db = db
self.client = client or get_mediakit_client()
self._cosyvoice = cosyvoice_service
self._voice_clone_repo = voice_clone_repo
def _get_cosyvoice(self):
"""延迟获取 CosyVoiceService(与 tts 路由一致,含 OSS 预签名配置)."""
if self._cosyvoice is None:
from app.dependencies import get_cosyvoice_service
self._cosyvoice = get_cosyvoice_service()
return self._cosyvoice
def _resolve_voice_id(self, voice_id: str, user_id: str) -> str:
"""将克隆音色 profile UUID 解析为 CosyVoice voice_id。
与 /tts/synthesize 保持一致:命中 profile → 校验归属 → 返回其 voice_id;
未命中(预置音色 ID 或克隆 CosyVoice voice_id)原样返回。
"""
if not voice_id:
return ""
if self._voice_clone_repo is None:
try:
from app.dependencies import get_voice_clone_profile_repository
self._voice_clone_repo = get_voice_clone_profile_repository(self.db)
except Exception:
return voice_id
try:
profile = self._voice_clone_repo.get(voice_id)
except Exception:
return voice_id
if profile is None:
return voice_id
if getattr(profile, "user_id", "") != user_id:
raise MediaKitError("无权访问该音色", code="VoiceForbidden")
if not getattr(profile, "voice_id", ""):
raise MediaKitError("音色克隆尚未完成,请稍后再试", code="VoiceNotReady")
return profile.voice_id
def _synthesize_and_persist_audio(
self,
*,
user_id: str,
job_id: str,
voice_id: str,
script_text: str,
speed: float,
emotion: str,
) -> str:
"""TTS 直生:调 CosyVoice 合成音频并转存 OSS,返回可公网访问的音频 URL.
Raises:
MediaKitError: 合成失败
"""
actual_voice_id = self._resolve_voice_id(voice_id, user_id)
cosyvoice = self._get_cosyvoice()
try:
result = cosyvoice.submit_synthesize_task(
text=script_text,
voice_id=actual_voice_id,
speed=speed,
emotion=normalize_emotion(emotion),
)
except CosyVoiceError as exc:
raise MediaKitError(f"TTS 合成失败: {exc}", code="TTSSynthesisFailed") from exc
except ValueError as exc:
raise MediaKitError(f"TTS 参数错误: {exc}", code="TTSInvalidParam") from exc
temp_url = result.get("audio_url", "")
if not temp_url:
raise MediaKitError("TTS 未返回音频 URL", code="TTSNoAudio")
# 转存到自家 OSS,避免临时 URL 过期导致 MediaKit 拉取失败
try:
audio_data = safe_download_bytes(
temp_url,
purpose="lipsync_tts_audio",
allowed_mime_types=ALLOWED_AUDIO_MIME_TYPES,
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("对口型 TTS 音频已转存 OSS: job_id=%s key=%s", job_id, storage_key)
return permanent_url
except Exception as exc:
logger.warning("TTS 音频转存 OSS 失败,回退临时 URL: job_id=%s err=%s", job_id, exc)
return temp_url
# ── 创建任务 ──────────────────────────────────────────────────────────
def create_job(
self,
*,
user_id: str,
video_url: str,
audio_url: str = "",
voice_id: str = "",
script_text: str = "",
speed: float = 1.0,
emotion: str = "",
enable_video_loop: bool = False,
project_id: str = "",
) -> LipsyncJobModel:
"""创建对口型任务.
两种输入模式:
- TTS 直生:voice_id + script_text(audio_url 留空)
→ 先创建 DB 记录(状态 tts_processing),再 dispatch Celery 异步任务
执行 TTS 合成 + MediaKit 提交。API 响应 <1s。
- 直接音频:提供 audio_url
→ 同步提交 MediaKit,状态直接设为 submitted。
Raises:
MediaKitError: 参数校验失败或 MediaKit 提交失败(仅直接音频模式)
"""
# 0. 输入校验
if not audio_url:
if not (voice_id and script_text):
raise MediaKitError(
"必须提供 audio_url 或 voice_id+script_text",
code="InvalidInput",
)
# 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,
project_id=project_id,
video_url=video_url,
audio_url=audio_url,
enable_video_loop=enable_video_loop,
voice_id=voice_id or "",
script_text=script_text or "",
speed=speed,
emotion=normalize_emotion(emotion),
status="tts_processing" if is_tts_mode else "pending",
)
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:
tts_synthesize_and_submit.apply_async(
args=(
job_id,
user_id,
voice_id,
script_text,
speed,
normalize_emotion(emotion),
)
)
except Exception as exc:
# 投递失败时立即把 job 标成 failed 并写入 error_message,
# 前端轮询时能直接看到失败原因,不会无限卡在 tts_processing。
logger.exception(
"Celery 任务提交失败,TTS 任务已创建但未触发执行: job_id=%s err=%s",
job_id,
exc,
)
job.status = "failed"
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)
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)
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
return job
# ── 查询任务 ──────────────────────────────────────────────────────────
def get_job(self, job_id: str, user_id: str) -> Optional[LipsyncJobModel]:
"""获取任务详情(用户隔离)."""
return (
self.db.query(LipsyncJobModel)
.filter(LipsyncJobModel.id == job_id, LipsyncJobModel.user_id == user_id)
.first()
)
def list_jobs(
self,
*,
user_id: str,
project_id: str = "",
status: str = "",
offset: int = 0,
limit: int = 20,
) -> tuple[list[LipsyncJobModel], int]:
"""获取任务列表(分页 + 用户隔离)."""
query = self.db.query(LipsyncJobModel).filter(LipsyncJobModel.user_id == user_id)
if project_id:
query = query.filter(LipsyncJobModel.project_id == project_id)
if status:
query = query.filter(LipsyncJobModel.status == status)
total = query.count()
items = query.order_by(LipsyncJobModel.created_at.desc()).offset(offset).limit(limit).all()
return items, total
# ── 更新任务状态(轮询) ──────────────────────────────────────────────
def refresh_job_status(self, job_id: str, user_id: str) -> Optional[LipsyncJobModel]:
"""从 MediaKit 拉取最新状态并更新本地记录.
Returns:
更新后的 Job,或 None(任务不存在/不属于该用户)
"""
job = self.get_job(job_id, user_id)
if job is None:
return None
# 终态不需要再轮询
if job.status in (STATUS_COMPLETED, "failed"):
return job
# 未提交的任务不轮询
if not job.mediakit_task_id:
return job
try:
status_data = self.client.get_task_status(job.mediakit_task_id)
except MediaKitError as exc:
logger.error("轮询对口型任务状态失败 [%s]: %s", job_id, exc)
return job
mk_status = status_data.get("status", STATUS_RUNNING)
logger.info("MediaKit 对口型状态 [%s]: %s", job_id, mk_status)
if mk_status == STATUS_COMPLETED:
result = status_data.get("result", {})
job.status = STATUS_COMPLETED
temp_url = result.get("video_url", "")
# 先以临时 URL 立即返回前端(前端可立即播放),再异步 Celery 任务转存自家 OSS(步骤⑦)
job.output_video_url = temp_url
job.output_duration = result.get("duration", 0.0)
job.completed_at = datetime.now(timezone.utc)
job.updated_at = datetime.now(timezone.utc)
self.db.commit()
# 异步转存到自家 OSS(注意:必须在 commit 之后 dispatch,避免 commit 失败任务已发出)
try:
from app.tasks.lipsync_tts import persist_output_video_task
persist_output_video_task.apply_async(args=(job_id, user_id, temp_url))
except Exception as exc:
logger.warning(
"提交输出视频异步转存任务失败,保留临时 URL: job_id=%s err=%s",
job_id,
exc,
)
self.db.refresh(job)
return job
elif mk_status == STATUS_FAILED:
error = status_data.get("error", {})
job.status = "failed"
job.error_message = error.get("message", "任务执行失败")
job.error_code = error.get("code", "TaskFailed")
job.completed_at = datetime.now(timezone.utc)
else:
# 中间状态(running/processing/queued 等)同步到 DB,避免前端永远卡在 submitted
if isinstance(mk_status, str) and mk_status:
job.status = mk_status
job.updated_at = datetime.now(timezone.utc)
self.db.commit()
self.db.refresh(job)
return job
def _persist_output_video(self, temp_url: str, job_id: str, user_id: str) -> str:
"""将 MediaKit 输出的临时视频 URL 转存到自家 OSS.
失败时回退返回原始临时 URL,不影响任务完成。
"""
if not temp_url:
return ""
try:
import httpx
with httpx.Client(timeout=180.0, follow_redirects=True) as client:
resp = client.get(temp_url)
resp.raise_for_status()
data = resp.content
storage = get_shared_storage_service()
storage_key = f"lipsync-outputs/{user_id}/{job_id}.mp4"
permanent_url = storage.upload_file(io.BytesIO(data), storage_key, content_type="video/mp4")
logger.info("对口型输出视频已转存 OSS: job_id=%s key=%s", job_id, storage_key)
return self._sign_media_url(permanent_url) or temp_url
except Exception as exc:
logger.warning("对口型输出视频转存 OSS 失败,回退临时 URL: job_id=%s err=%s", job_id, exc)
return temp_url
def _sign_media_url(self, url: str) -> str:
"""对自家 OSS 私有桶 URL 重签长有效期预签名,供 MediaKit 拉取 / 前端播放。
- 裸 public_url(upload_file 返回,不带签名)→ 私有桶匿名访问 403,重签。
- 已带签名但即将过期的 URL(如前端 1h 预签名)→ 抽 storage_key 后重签。
- 外部 URL(CosyVoice/MediaKit 临时链接,非本桶 host)→ 原样透传。
- 任何异常都降级原样返回,不阻断主流程。
"""
if not url:
return url
try:
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 # 非自家 OSS(外部临时链接),不处理
signed = storage.get_download_url(url, expires_seconds=MEDIAKIT_URL_TTL_SECONDS)
return signed or url
except Exception as exc: # noqa: BLE001 - 签名失败不阻断,降级原 URL
logger.warning("对口型 URL 重签失败,原样返回: url_prefix=%s err=%s", url[:80], exc)
return url
# ── 取消任务 ──────────────────────────────────────────────────────────
def cancel_job(self, job_id: str, user_id: str) -> Optional[LipsyncJobModel]:
"""取消任务(仅 pending/tts_processing/submitted 状态可取消)."""
job = self.get_job(job_id, user_id)
if job is None:
return None
if job.status in ("pending", "tts_processing", "submitted"):
job.status = "cancelled"
job.updated_at = datetime.now(timezone.utc)
self.db.commit()
self.db.refresh(job)
return job