146effcf07
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 1s
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / PR Build Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / ACR Image Cleanup (pull_request) Has been skipped
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 1m31s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 1m55s
CI/CD Pipeline / Unit Tests (pull_request) Failing after 2m8s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m3s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 3m8s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 3m18s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 3m17s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 5m17s
AI Code Review / AI Code Review (pull_request) Successful in 6m30s
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
CI/CD Pipeline / Validate - Style (pull_request) Has been cancelled
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been cancelled
根因: - 声音克隆走 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日志,便于日志排查
118 lines
4.5 KiB
Python
Executable File
118 lines
4.5 KiB
Python
Executable File
"""Voice clone tasks - process voice clone requests via CosyVoice API."""
|
||
|
||
import logging
|
||
|
||
from celery import Task
|
||
from celery.exceptions import Retry
|
||
from video_processing.oss_helpers import get_signed_download_url
|
||
from worker_app.celery_app import celery_app
|
||
from worker_app.db import SessionLocal
|
||
|
||
from packages.adapters.sqlalchemy_impl.voice_clone_profile_repository import (
|
||
SQLAlchemyVoiceCloneProfileRepository,
|
||
)
|
||
from packages.application.cosyvoice_service import (
|
||
CosyVoiceError,
|
||
CosyVoiceService,
|
||
CosyVoiceTimeoutError,
|
||
)
|
||
from packages.application.voice_clone.workflow import VoiceCloneWorkflowService
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
|
||
@celery_app.task(bind=True, max_retries=2, name="worker.process_voice_clone")
|
||
def process_voice_clone(self: Task, profile_id: str) -> dict:
|
||
"""处理音色克隆任务。
|
||
|
||
通过 VoiceCloneWorkflowService.poll_and_process_clone() 轮询 CosyVoice
|
||
克隆任务状态,更新 VoiceCloneProfile。
|
||
|
||
重试策略(Celery 5.x bind=True 模式):
|
||
- max_retries=2:最多重试 2 次(共执行 3 次),超过后抛出 MaxRetriesExceededError
|
||
- CosyVoiceTimeoutError:网络超时属于临时性故障,使用 countdown=30 延迟 30 秒后重试
|
||
- CosyVoiceError:API 业务错误(如任务失败),属于永久性故障,不重试直接标记 failed
|
||
- Exception:未知错误,不重试直接标记 failed,避免无限重试掩盖 bug
|
||
- Retry 异常:Celery 内部重试信号,必须向上传播不能被捕获
|
||
|
||
Args:
|
||
profile_id: VoiceCloneProfile ID
|
||
|
||
Returns:
|
||
dict: {"ok": True, "profile_id": str, "voice_id": str}
|
||
"""
|
||
# P2-2 修复:session 初始化为 None,避免 SessionLocal() 抛异常时
|
||
# finally 块中 session.close() 触发 UnboundLocalError
|
||
session = None
|
||
logger.info(f"Voice clone task started: profile_id={profile_id}")
|
||
try:
|
||
session = SessionLocal()
|
||
repo = SQLAlchemyVoiceCloneProfileRepository(session)
|
||
workflow = VoiceCloneWorkflowService(
|
||
repository=repo,
|
||
cosyvoice_service=CosyVoiceService(
|
||
audio_url_signer=lambda url: get_signed_download_url(url, expires_seconds=86400) or url
|
||
),
|
||
)
|
||
|
||
updated_profile = workflow.poll_and_process_clone(profile_id, timeout=300)
|
||
session.commit()
|
||
|
||
logger.info(f"Voice clone completed: profile_id={profile_id}, " f"voice_id={updated_profile.voice_id}")
|
||
return {
|
||
"ok": True,
|
||
"profile_id": profile_id,
|
||
"voice_id": updated_profile.voice_id,
|
||
}
|
||
|
||
except Retry:
|
||
# Celery Retry 异常必须向上传播,不能被后续 except 捕获
|
||
raise
|
||
|
||
except CosyVoiceTimeoutError as e:
|
||
logger.error(f"Voice clone timeout for {profile_id}: {e}")
|
||
if session is not None:
|
||
session.rollback()
|
||
# 超时属于临时性故障,延迟 30 秒后重试
|
||
raise self.retry(exc=e, countdown=30) from e
|
||
|
||
except CosyVoiceError as e:
|
||
logger.error(f"Voice clone failed for {profile_id}: {e}")
|
||
if session is not None:
|
||
session.rollback()
|
||
# API 业务错误属于永久性故障,不重试,标记 profile 为 failed
|
||
try:
|
||
if session is not None:
|
||
profile = repo.get(profile_id)
|
||
if profile is not None:
|
||
profile.mark_failed(str(e))
|
||
repo.update(profile)
|
||
session.commit()
|
||
except Exception as inner_e:
|
||
logger.error(f"Failed to mark profile as failed: {inner_e}")
|
||
if session is not None:
|
||
session.rollback()
|
||
return {"ok": False, "profile_id": profile_id, "error": str(e)}
|
||
|
||
except Exception as e:
|
||
logger.error(f"Voice clone unexpected error for {profile_id}: {e}")
|
||
if session is not None:
|
||
session.rollback()
|
||
# 未知错误不重试,标记 profile 为 failed,避免无限重试掩盖 bug
|
||
try:
|
||
if session is not None:
|
||
profile = repo.get(profile_id)
|
||
if profile is not None:
|
||
profile.mark_failed(str(e))
|
||
repo.update(profile)
|
||
session.commit()
|
||
except Exception as inner_e:
|
||
logger.error(f"Failed to mark profile as failed: {inner_e}")
|
||
if session is not None:
|
||
session.rollback()
|
||
return {"ok": False, "profile_id": profile_id, "error": str(e)}
|
||
|
||
finally:
|
||
if session is not None:
|
||
session.close()
|