Files
xiaoxia-saas/packages/adapters/sqlalchemy_impl/voice_clone_profile_repository.py
T
CI Bot a33532113b fix(voice-clone): 处理中卡死兜底 — worker重启/Celery消息丢失后processing记录永久卡住问题
根因:
- 声音克隆走 API提交CosyVoice→Celery worker轮询5分钟→mark ready/failed
- 当worker容器重启、进程OOM或Celeryprefetch消息丢失时,已进入processing状态的voice_clone_profile没有兜底恢复机制
- 对比ingest链路有cleanup_stale_ingest_jobs beat任务+worker_ready启动恢复,voice_clone链路缺失
- 结果: 用户看到"克隆处理中,请稍候..."永久转圈,刷新也不变

修复:
1. SQLAlchemyVoiceCloneProfileRepository新增cleanup_stale_processing方法:
   扫描updated_at超过10分钟的processing记录(正常克隆<5分钟),标记failed并提示用户重试
2. _startup.py新增recover_stale_voice_clones_on_startup+worker_ready信号:
   worker启动时一次性扫描恢复,用户重启/部署后立即解锁卡死任务
3. cleanup.py新增scheduled_cleanup_stale_voice_clones beat任务:
   每5分钟巡检兜底,防止运行期OOM/消息丢失导致的新卡死
4. celery_app.py beat_schedule注册新任务
5. voice_clone.py任务入口加received日志,便于日志排查
2026-09-26 14:51:57 +08:00

192 lines
7.1 KiB
Python

"""SQLAlchemy implementation of VoiceCloneProfileRepository."""
from __future__ import annotations
from typing import Optional
from sqlalchemy.orm import Session
from packages.adapters.sqlalchemy_impl.models import VoiceCloneProfileModel
from packages.domain.voice_clone_profile import VoiceCloneProfile, VoiceCloneStatus
class SQLAlchemyVoiceCloneProfileRepository:
"""SQLAlchemy 音色克隆档案仓储。"""
def __init__(self, session: Session) -> None:
self.session = session
def create(self, profile: VoiceCloneProfile) -> VoiceCloneProfile:
model = VoiceCloneProfileModel(
id=profile.id,
user_id=profile.user_id,
name=profile.name,
description=profile.description,
source_audio_url=profile.source_audio_url,
voice_id=profile.voice_id,
voice_model=profile.voice_model,
language=profile.language,
gender=profile.gender,
status=profile.status,
error_message=profile.error_message,
retry_count=profile.retry_count,
max_retries=profile.max_retries,
metadata_=profile.metadata,
)
self.session.add(model)
self.session.commit()
self.session.refresh(model)
return self._model_to_entity(model)
def get(self, profile_id: str) -> Optional[VoiceCloneProfile]:
model = (
self.session.query(VoiceCloneProfileModel)
.filter(
VoiceCloneProfileModel.id == profile_id,
VoiceCloneProfileModel.status != "deleted",
)
.first()
)
if model is None:
return None
return self._model_to_entity(model)
def update(self, profile: VoiceCloneProfile) -> VoiceCloneProfile:
model = self.session.query(VoiceCloneProfileModel).filter(VoiceCloneProfileModel.id == profile.id).first()
if model is None:
raise ValueError(f"VoiceCloneProfile {profile.id} not found")
model.name = profile.name
model.description = profile.description
model.source_audio_url = profile.source_audio_url
model.voice_id = profile.voice_id
model.voice_model = profile.voice_model
model.language = profile.language
model.gender = profile.gender
model.status = profile.status
model.error_message = profile.error_message
model.retry_count = profile.retry_count
model.max_retries = profile.max_retries
model.metadata_ = profile.metadata
self.session.commit()
self.session.refresh(model)
return self._model_to_entity(model)
def delete(self, profile_id: str) -> bool:
model = self.session.query(VoiceCloneProfileModel).filter(VoiceCloneProfileModel.id == profile_id).first()
if model is None:
return False
model.status = "deleted"
self.session.commit()
return True
def list_by_user(
self,
user_id: str,
*,
status: Optional[str] = None,
limit: int = 50,
offset: int = 0,
) -> list[VoiceCloneProfile]:
query = self.session.query(VoiceCloneProfileModel).filter(
VoiceCloneProfileModel.user_id == user_id,
VoiceCloneProfileModel.status != "deleted",
)
if status:
query = query.filter(VoiceCloneProfileModel.status == status)
query = query.order_by(VoiceCloneProfileModel.created_at.desc())
models = query.offset(offset).limit(limit).all()
return [self._model_to_entity(m) for m in models]
def count_by_user(self, user_id: str, *, status: Optional[str] = None) -> int:
query = self.session.query(VoiceCloneProfileModel).filter(
VoiceCloneProfileModel.user_id == user_id,
VoiceCloneProfileModel.status != "deleted",
)
if status:
query = query.filter(VoiceCloneProfileModel.status == status)
return query.count()
def find_by_voice_id(self, voice_id: str) -> Optional[VoiceCloneProfile]:
model = (
self.session.query(VoiceCloneProfileModel)
.filter(
VoiceCloneProfileModel.voice_id == voice_id,
VoiceCloneProfileModel.status != "deleted",
)
.first()
)
if model is None:
return None
return self._model_to_entity(model)
def find_profile_ids_by_voice_ids(self, voice_ids: list[str]) -> dict[str, str]:
"""批量查询 voice_id → profile_id 映射。用于填充统一列表的 voice_clone_profile_id。"""
if not voice_ids:
return {}
rows = (
self.session.query(
VoiceCloneProfileModel.voice_id,
VoiceCloneProfileModel.id,
)
.filter(
VoiceCloneProfileModel.voice_id.in_(voice_ids),
VoiceCloneProfileModel.status != "deleted",
)
.all()
)
return {voice_id: profile_id for voice_id, profile_id in rows}
def cleanup_stale_processing(self, timeout_minutes: int = 10) -> int:
"""清理超时卡在 processing 的克隆档案。
worker 重启、Celery 任务丢失或 OOM 被杀时,processing 档案会永久卡住。
updated_at < NOW() - timeout_minutes 的 processing 记录,标记为 failed
并附带明确错误信息,用户可在前端点击「重试」。
Args:
timeout_minutes: 超时分钟数,默认 10 分钟(正常克隆 < 5 分钟)
Returns:
清理的记录数
"""
from datetime import datetime, timedelta, UTC
cutoff = datetime.now(UTC) - timedelta(minutes=timeout_minutes)
models = (
self.session.query(VoiceCloneProfileModel)
.filter(
VoiceCloneProfileModel.status == "processing",
VoiceCloneProfileModel.updated_at < cutoff,
)
.all()
)
count = 0
for model in models:
model.status = "failed"
model.error_message = f"克隆任务执行超时(超过 {timeout_minutes} 分钟未更新,可能因服务重启中断),请重试"
count += 1
if count > 0:
self.session.commit()
return count
@staticmethod
def _model_to_entity(model: VoiceCloneProfileModel) -> VoiceCloneProfile:
return VoiceCloneProfile(
id=model.id,
user_id=model.user_id,
name=model.name,
description=model.description or "",
source_audio_url=model.source_audio_url or "",
voice_id=model.voice_id or "",
voice_model=model.voice_model or "",
language=model.language or "zh-CN",
gender=model.gender or "unknown",
status=VoiceCloneStatus(model.status) if model.status else VoiceCloneStatus.PENDING,
error_message=model.error_message or "",
retry_count=model.retry_count or 0,
max_retries=model.max_retries or 3,
metadata=model.metadata_ or {},
created_at=model.created_at,
updated_at=model.updated_at,
)