"""叙事剪辑前置服务 — #1970 PR3. 叙事模式(assembly_mode='narrative')在生成任务入队前同步完成: 1. 按 script_id 读取文案(归属校验); 2. 按 tts_voice_source 解析音色(preset=CosyVoice 音色 id;clone=克隆档案 id, 解析档案归属并取其 CosyVoice voice_id); 3. 同步 TTS 合成(复用 tts_job 现有 workflow:提交即同步返回,未完成则轮询兜底), 失败直接抛 NarrativeError(HTTP 层转 4xx,任务不入队); 4. 把合成音频转存为配音库 audio asset(与 /tts/jobs/{id}/save-to-library 同一套 存储路径与元信息约定),返回 asset_id —— 下游仍以 voice_library_id(实为 audio asset id)消费,渲染链路零改动。 """ from __future__ import annotations import json import logging import subprocess import tempfile from dataclasses import dataclass from pathlib import Path from typing import Any from sqlalchemy.orm import Session from packages.adapters.sqlalchemy_impl.models import ScriptModel from packages.application.cosyvoice_service import CosyVoiceService from packages.application.tts_job.use_cases import CreateTTSJobUseCase from packages.application.tts_job.workflow import TTSWorkflowService from packages.domain import Asset, AssetLibrary, AssetLibraryKind, AssetStatus, ClassificationStatus from packages.shared.storage import SharedStorageService logger = logging.getLogger(__name__) _SYNTH_TIMEOUT = 180.0 # 叙事配音在 HTTP 请求内同步等待,长文案分段合成时留出余量 _CONTENT_TYPE_MAP = {"mp3": "audio/mpeg", "wav": "audio/wav", "pcm": "audio/pcm", "opus": "audio/opus"} class NarrativeError(Exception): """叙事模式前置处理失败(文案/音色/TTS/落库)。""" def __init__(self, message: str, *, status_code: int = 400) -> None: super().__init__(message) self.message = message self.status_code = status_code @dataclass(slots=True) class NarrativeContext: """叙事模式前置处理结果。""" script: ScriptModel voice_asset_id: str tts_job_id: str audio_duration: float def _find_or_create_voice_library( *, user_id: str, project_repository: Any, asset_library_repository: Any, ) -> AssetLibrary: """找到(或自动创建)用户 voice 素材库;与 tts.py 保存配音库逻辑一致。""" projects = project_repository.find_accessible_projects(user_id) if not projects: raise NarrativeError("没有可用的项目,无法保存叙事配音", status_code=400) for project in projects: for lib in asset_library_repository.find_by_project(project.id): kind = lib.kind.value if hasattr(lib.kind, "value") else lib.kind if kind == AssetLibraryKind.VOICE.value: return lib project = projects[0] library = AssetLibrary.create(project_id=project.id, name="配音素材库", kind=AssetLibraryKind.VOICE) from sqlalchemy.exc import IntegrityError try: return asset_library_repository.create(library) except IntegrityError: session = getattr(asset_library_repository, "session", None) if session is not None: try: session.rollback() except Exception: # noqa: BLE001 - 回滚失败不影响重查 logger.warning("IntegrityError 后回滚 session 失败", exc_info=True) for lib in asset_library_repository.find_by_project(project.id): kind = lib.kind.value if hasattr(lib.kind, "value") else lib.kind if kind == AssetLibraryKind.VOICE.value: return lib raise NarrativeError("配音素材库创建失败,请重试", status_code=500) from None def _resolve_voice( *, user_id: str, tts_voice_id: str, tts_voice_source: str, voice_clone_repository: Any, ) -> tuple[str, str]: """解析音色 → (CosyVoice voice_id, voice_clone_profile_id)。""" if tts_voice_source == "clone": profile = voice_clone_repository.get(tts_voice_id) if profile is None: raise NarrativeError("克隆音色不存在", status_code=404) if profile.user_id != user_id: raise NarrativeError("无权使用该克隆音色", status_code=403) if not profile.voice_id: raise NarrativeError("音色克隆尚未完成,请稍后再试", status_code=400) return profile.voice_id, profile.id # preset:tts_voice_id 即 CosyVoice 音色 id;与 /tts 端点一致, # 若前端误传克隆档案 UUID,同样兼容解析。 profile = voice_clone_repository.get(tts_voice_id) if profile is not None: if profile.user_id != user_id: raise NarrativeError("无权使用该音色", status_code=403) if not profile.voice_id: raise NarrativeError("音色克隆尚未完成,请稍后再试", status_code=400) return profile.voice_id, profile.id return tts_voice_id, "" def _save_tts_job_as_voice_asset( *, job: Any, user_id: str, name: str, project_repository: Any, asset_library_repository: Any, asset_repository: Any, storage_service: SharedStorageService, ) -> Asset: """把已完成 TTS job 的音频转存为配音库 audio asset(同 save-to-library 约定)。""" if not job.output_audio_url and not job.output_audio_key: raise NarrativeError("TTS 合成缺少输出音频", status_code=502) library = _find_or_create_voice_library( user_id=user_id, project_repository=project_repository, asset_library_repository=asset_library_repository, ) audio_format = (job.format or "mp3").strip() or "mp3" content_type = _CONTENT_TYPE_MAP.get(audio_format, "audio/mpeg") storage_key = f"uploads/voice/tts/{job.id}.{audio_format}" tmp_path: Path | None = None audio_duration: float | None = None file_size = 0 try: with tempfile.NamedTemporaryFile(suffix=f".{audio_format}", delete=False) as tmp: tmp_path = Path(tmp.name) download_source = job.output_audio_key or job.output_audio_url downloaded = storage_service.download_asset(download_source, tmp_path) if not downloaded or not tmp_path.exists() or tmp_path.stat().st_size == 0: raise NarrativeError("叙事配音音频转存失败", status_code=502) file_size = tmp_path.stat().st_size storage_service.upload_file(tmp_path, storage_key, content_type=content_type) try: proc = subprocess.run( [ "ffprobe", "-v", "quiet", "-print_format", "json", "-show_format", str(tmp_path), ], capture_output=True, text=True, timeout=10, ) if proc.returncode == 0: dur = float(json.loads(proc.stdout).get("format", {}).get("duration", 0)) if dur > 0: audio_duration = dur except Exception: # noqa: BLE001 - ffprobe 仅用于时长兜底 logger.warning("叙事配音 ffprobe 时长提取失败: job_id=%s", job.id, exc_info=True) except NarrativeError: raise except Exception as e: # noqa: BLE001 logger.error("叙事配音转存失败: job_id=%s, error=%s", job.id, e, exc_info=True) raise NarrativeError("叙事配音音频转存失败", status_code=502) from e finally: if tmp_path and tmp_path.exists(): try: tmp_path.unlink() except OSError: pass metadata_: dict[str, object] = { "source": "tts_job", "tts_job_id": job.id, "narrative": True, "format": job.format, "sample_rate": job.sample_rate, "voice_id": job.voice_id, "voice_name": job.voice_model or "", } if job.metadata: for key in ("speed", "language"): if key in job.metadata: metadata_[key] = job.metadata[key] asset = Asset.create( project_id=library.project_id, library_id=library.id, name=name or f"叙事配音-{job.id[:8]}", storage_key=storage_key, mime_type=content_type, metadata=metadata_, file_size=file_size, duration=job.duration or audio_duration or None, status=AssetStatus.READY, classification_status=ClassificationStatus.PENDING, uploaded_by_user_id=user_id, ) try: return asset_repository.create(asset) except Exception as e: # noqa: BLE001 logger.error("叙事配音 asset 落库失败,清理 OSS: %s, error=%s", storage_key, e, exc_info=True) try: storage_service.delete_file(storage_key) except Exception: # noqa: BLE001 logger.warning("清理孤儿 OSS 文件失败: %s", storage_key, exc_info=True) raise NarrativeError("叙事配音保存失败,请重试", status_code=502) from e def prepare_narrative_voice( *, db: Session, user_id: str, script_id: str, tts_voice_id: str, tts_voice_source: str, tts_repository: Any, cosyvoice_service: CosyVoiceService, voice_clone_repository: Any, asset_repository: Any, asset_library_repository: Any, project_repository: Any, storage_service: SharedStorageService, points_enabled: bool = False, is_member: bool = False, member_type: str | None = None, ) -> NarrativeContext: """叙事模式入队前同步合成配音并落为 audio asset。 Raises: NarrativeError: 文案缺失/归属不符、音色不可用、TTS 失败、转存失败。 """ script = db.query(ScriptModel).filter(ScriptModel.id == script_id, ScriptModel.user_id == user_id).first() if script is None: raise NarrativeError("文案不存在或无权使用", status_code=404) content = (script.content or "").strip() if not content: raise NarrativeError("文案内容为空,无法合成配音", status_code=400) actual_voice_id, clone_profile_id = _resolve_voice( user_id=user_id, tts_voice_id=tts_voice_id, tts_voice_source=tts_voice_source, voice_clone_repository=voice_clone_repository, ) use_case = CreateTTSJobUseCase(tts_repository) job = use_case.execute( user_id=user_id, input_text=content, voice_id=actual_voice_id, voice_clone_profile_id=clone_profile_id, metadata={"speed": 1.0, "emotion": "", "language": "zh-CN", "narrative": True, "script_id": script_id}, ) workflow = TTSWorkflowService(repository=tts_repository, cosyvoice_service=cosyvoice_service) try: job = workflow.start_synthesis(job.id) if not job.is_completed: job = workflow.poll_and_process_synthesis(job.id, timeout=_SYNTH_TIMEOUT) except Exception as e: # noqa: BLE001 - 同步合成异常统一转 NarrativeError logger.error("叙事配音 TTS 合成失败: job_id=%s, error=%s", job.id, e, exc_info=True) try: workflow.process_synthesis_failure(job.id, str(e)) except Exception: # noqa: BLE001 logger.warning("标记叙事 TTS job 失败出错: job_id=%s", job.id, exc_info=True) raise NarrativeError(f"配音合成失败:{e}", status_code=502) from e if not job.is_completed: raise NarrativeError("配音合成未完成,请稍后重试", status_code=504) asset = _save_tts_job_as_voice_asset( job=job, user_id=user_id, name=(script.title or "叙事配音")[:60], project_repository=project_repository, asset_library_repository=asset_library_repository, asset_repository=asset_repository, storage_service=storage_service, ) return NarrativeContext( script=script, voice_asset_id=asset.id, tts_job_id=job.id, audio_duration=float(job.duration or asset.duration or 0.0), )