"""AI数字人渲染合成 Service — #1798. 职责: - 创建/查询/取消渲染任务 - 调用 Celery 异步任务执行渲染 - B-roll 合成 + 标题叠加 + 封面提取 - 用户隔离 """ from __future__ import annotations import base64 import binascii import logging import os import subprocess import tempfile import uuid from datetime import UTC, datetime from typing import Any, Optional from sqlalchemy.orm import Session from packages.adapters.sqlalchemy_impl.models import ( AiAvatarRenderJob, LipsyncJobModel, ScriptModel, ) from packages.domain.video_filter_builder import ( build_broll_overlay_filter, build_title_drawtext_filter, build_title_overlay_filter, ) from packages.shared.storage import get_shared_storage_service logger = logging.getLogger(__name__) class AiAvatarRenderError(Exception): """渲染服务异常.""" def __init__(self, message: str, code: str = "RenderError"): self.code = code super().__init__(message) class AiAvatarRenderService: """AI数字人渲染合成 Service.""" def __init__(self, db: Session): self.db = db # ── 创建任务 ────────────────────────────────────────────────────────── def create_render_job( self, *, user_id: str, lipsync_job_id: str, script_id: str = "", b_roll_segments: list[dict[str, Any]] | None = None, title_config: dict[str, Any], cover_config: dict[str, Any], project_id: str = "", ) -> AiAvatarRenderJob: """创建渲染任务. Raises: AiAvatarRenderError: 校验失败 """ # 1. 验证对口型任务 lipsync_job = ( self.db.query(LipsyncJobModel) .filter( LipsyncJobModel.id == lipsync_job_id, LipsyncJobModel.user_id == user_id, ) .first() ) if lipsync_job is None: raise AiAvatarRenderError("对口型任务不存在", code="LipsyncJobNotFound") if lipsync_job.status != "completed": raise AiAvatarRenderError( f"对口型任务状态为 {lipsync_job.status},仅 completed 状态可渲染", code="LipsyncJobNotCompleted", ) if not lipsync_job.output_video_url: raise AiAvatarRenderError("对口型任务输出视频 URL 为空", code="LipsyncJobNoOutput") # 2. 验证文案归属(仅当选了文案库条目时;手动输入文案直生场景 script_id 可空) script_id = (script_id or "").strip() if script_id: script = ( self.db.query(ScriptModel) .filter( ScriptModel.id == script_id, ScriptModel.user_id == user_id, ) .first() ) if script is None: raise AiAvatarRenderError("文案不存在或无权访问", code="ScriptNotFound") # 3. 创建渲染任务 job_id = str(uuid.uuid4()) job = AiAvatarRenderJob( id=job_id, user_id=user_id, project_id=project_id, lipsync_job_id=lipsync_job_id, script_id=script_id, b_roll_segments=[s if isinstance(s, dict) else s.model_dump() for s in (b_roll_segments or [])], title_config=title_config, cover_config=cover_config, status="pending", ) self.db.add(job) self.db.flush() job.submitted_at = datetime.now(UTC) self.db.commit() self.db.refresh(job) return job # ── 查询任务 ────────────────────────────────────────────────────────── def get_render_job(self, job_id: str, user_id: str) -> Optional[AiAvatarRenderJob]: """获取渲染任务详情(用户隔离).""" return ( self.db.query(AiAvatarRenderJob) .filter( AiAvatarRenderJob.id == job_id, AiAvatarRenderJob.user_id == user_id, ) .first() ) def list_render_jobs( self, *, user_id: str, project_id: str = "", status: str = "", offset: int = 0, limit: int = 20, ) -> tuple[list[AiAvatarRenderJob], int]: """获取渲染任务列表(分页 + 用户隔离).""" query = self.db.query(AiAvatarRenderJob).filter(AiAvatarRenderJob.user_id == user_id) if project_id: query = query.filter(AiAvatarRenderJob.project_id == project_id) if status: query = query.filter(AiAvatarRenderJob.status == status) total = query.count() items = query.order_by(AiAvatarRenderJob.created_at.desc()).offset(offset).limit(limit).all() return items, total # ── 取消任务 ────────────────────────────────────────────────────────── def cancel_render_job(self, job_id: str, user_id: str) -> Optional[AiAvatarRenderJob]: """取消渲染任务(仅 pending 状态可取消).""" job = self.get_render_job(job_id, user_id) if job is None: return None if job.status in ("pending", "submitted"): job.status = "cancelled" job.updated_at = datetime.now(UTC) self.db.commit() self.db.refresh(job) return job # ── 重试任务 ────────────────────────────────────────────────────────── def retry_render_job(self, job_id: str, user_id: str) -> Optional[AiAvatarRenderJob]: """重试失败的渲染任务.""" job = self.get_render_job(job_id, user_id) if job is None: return None if job.status != "failed": return None job.status = "pending" job.progress = 0 job.error_message = "" job.output_video_url = "" job.output_cover_url = "" job.output_duration = 0.0 job.started_at = None job.completed_at = None job.updated_at = datetime.now(UTC) self.db.commit() self.db.refresh(job) return job # ── 执行渲染(Celery 异步调用) ────────────────────────────────────── def execute_render(self, job_id: str) -> None: """执行渲染管线. 由 Celery 异步任务调用,流程: 1. 下载对口型输出视频 (20%) 2. 构建 FFmpeg 滤镜链 (40%) 3. 执行 FFmpeg 渲染 (80%) 4. 上传到 OSS (95%) — 封面不再自动生成,改由前端主动抽帧 5. 更新任务状态 (100%) """ job = self.db.query(AiAvatarRenderJob).filter(AiAvatarRenderJob.id == job_id).first() if job is None: logger.error("渲染任务不存在: %s", job_id) return if job.status == "cancelled": logger.info("渲染任务已取消: %s", job_id) return try: # 更新状态为 processing job.status = "processing" job.started_at = datetime.now(UTC) job.progress = 5 job.updated_at = datetime.now(UTC) self.db.commit() # 获取对口型任务信息 lipsync_job = self.db.query(LipsyncJobModel).filter(LipsyncJobModel.id == job.lipsync_job_id).first() if lipsync_job is None: raise AiAvatarRenderError("关联的对口型任务不存在", code="LipsyncJobNotFound") # 1. 下载对口型输出视频 (20%) input_video_path = self._download_video(lipsync_job.output_video_url) job.progress = 20 self.db.commit() # 2. 构建 FFmpeg 滤镜链 (40%) # 用 ffprobe 探测输入视频分辨率,确保 B-roll 缩放与标题位置与实际输出一致。 # AI 数字人对口型输出为 9:16 竖屏,默认兜底 720x1280;探测失败时使用默认值不阻断渲染。 output_width, output_height = self._probe_video_resolution(input_video_path) if output_width <= 0 or output_height <= 0: output_width, output_height = 720, 1280 logger.info( "[数字人渲染] ffprobe 探测分辨率失败或无效,使用默认竖屏尺寸 %sx%s", output_width, output_height, ) else: logger.info("[数字人渲染] 探测输入视频分辨率: %sx%s", output_width, output_height) broll_filter, broll_label = build_broll_overlay_filter( b_roll_segments=job.b_roll_segments, video_duration=lipsync_job.output_duration, output_width=output_width, output_height=output_height, ) # 标题叠加路径:优先前端 Canvas 渲染的 PNG 图层(所见即所得), # 无 title_image_dataurl 时降级到 drawtext 重画文字。 title_cfg = job.title_config if isinstance(job.title_config, dict) else {} title_dataurl = (title_cfg or {}).get("title_image_dataurl") if title_cfg else None use_title_png = isinstance(title_dataurl, str) and title_dataurl.startswith("data:image/") title_input_index = 1 + len(job.b_roll_segments or []) if use_title_png else None job.progress = 40 self.db.commit() # 3. 执行 FFmpeg 渲染 (80%) with tempfile.TemporaryDirectory() as tmpdir: output_video_path = os.path.join(tmpdir, "output.mp4") # 在临时目录里解码保存标题 PNG(with 退出自动清理) title_png_path: Optional[str] = None extra_inputs: list[str] = [] title_filter = None if use_title_png: try: title_png_path = os.path.join(tmpdir, f"title_{job.id}.png") self._save_title_dataurl_to_file(title_dataurl, dst_path=title_png_path) extra_inputs.append(title_png_path) logger.info( "[数字人渲染] 标题 PNG 已保存: %s (input index %d)", title_png_path, title_input_index ) except Exception as exc: logger.warning("[数字人渲染] 标题 PNG 解码/保存失败,降级 drawtext: %s", exc) title_png_path = None extra_inputs = [] # 构建标题滤镜 final_label = None if title_png_path and title_input_index is not None: title_input_label = f"[{title_input_index}:v]" base_label = f"[{broll_label}]" if broll_label else "[0:v]" title_filter = build_title_overlay_filter( title_cfg, output_width=output_width, output_height=output_height, title_png_path=title_png_path, title_input_label=title_input_label, base_label=base_label, output_label="vout_titled", ) if not title_filter: # build 返回 None → 文件不存在(极端并发情况),降级 drawtext title_png_path = None extra_inputs = [] if title_png_path: # overlay 路径 if broll_filter and title_filter: filter_complex = broll_filter + f";{title_filter}" elif broll_filter: filter_complex = broll_filter final_label = broll_label elif title_filter: filter_complex = title_filter else: filter_complex = "" if title_filter: final_label = "vout_titled" elif not final_label: final_label = None else: # 降级:drawtext 重画文字 title_filter = build_title_drawtext_filter( title_cfg, output_width=output_width, output_height=output_height, ) if broll_filter and title_filter: filter_complex = broll_filter + f";[{broll_label}]{title_filter}[vout_titled]" final_label = "vout_titled" elif broll_filter: filter_complex = broll_filter final_label = broll_label elif title_filter: filter_complex = f"[0:v]{title_filter}[vout_titled]" final_label = "vout_titled" else: filter_complex = "" final_label = None cmd_list = self._build_ffmpeg_command( input_video=input_video_path, b_roll_segments=job.b_roll_segments, extra_inputs=extra_inputs, filter_complex=filter_complex, final_label=final_label, output_path=output_video_path, ) try: render_result = subprocess.run( cmd_list, capture_output=True, text=True, timeout=600, ) except subprocess.TimeoutExpired as exc: raise AiAvatarRenderError( "FFmpeg 渲染超时(600s)", code="FFmpegTimeout", ) from exc if render_result.returncode != 0: stderr_tail = (render_result.stderr or "").strip()[-800:] raise AiAvatarRenderError( f"FFmpeg 渲染失败,退出码: {render_result.returncode}, stderr: {stderr_tail}", code="FFmpegFailed", ) job.progress = 80 self.db.commit() # 4/5. 上传成片到 OSS (95%) —— 已砍掉自动抽封面逻辑(步骤⑤); # 封面由前端在渲染完成后通过 /smart-cover 接口主动从成片抽帧,不阻塞渲染链路。 output_video_url = self._upload_to_oss(output_video_path, f"ai-avatar/{job_id}/output.mp4") job.output_video_url = output_video_url # 封面透传:如果用户已在 cover_config 中选定封面 URL(mode=upload 的自定义上传 或 # mode=auto_frame 已有的智能封面结果),直接透传到 output_cover_url,不再重新截帧。 if isinstance(job.cover_config, dict): _pre_cover_url = ( job.cover_config.get("url") or job.cover_config.get("imageUrl") or job.cover_config.get("cover_url") or "" ) if _pre_cover_url: job.output_cover_url = _pre_cover_url logger.info("[数字人渲染] 使用用户已选定封面 URL: job_id=%s", job_id) # 获取输出视频时长 job.output_duration = lipsync_job.output_duration job.progress = 95 self.db.commit() # 6. 完成 job.status = "completed" job.progress = 100 job.completed_at = datetime.now(UTC) job.updated_at = datetime.now(UTC) self.db.commit() logger.info("渲染任务完成: %s", job_id) # 7. 渲染完成,停留在「待选封面」状态:不自动入库。 # 用户在前端选好封面、点「完成」后,由 /{job_id}/finalize 接口显式入库。 logger.info("渲染任务完成,等待用户选择封面后入库: job_id=%s", job_id) except AiAvatarRenderError as exc: job.status = "failed" job.error_message = str(exc) job.updated_at = datetime.now(UTC) self.db.commit() logger.error("渲染任务失败 [%s]: %s", job_id, exc) raise except Exception as exc: job.status = "failed" job.error_message = f"渲染异常: {str(exc)}" job.updated_at = datetime.now(UTC) self.db.commit() logger.exception("渲染任务异常 [%s]", job_id) raise def _persist_to_library(self, job: AiAvatarRenderJob, cover_url: Optional[str] = None): """将渲染结果写入成片库,返回 GeneratedVideo 领域对象. Args: job: 渲染任务(必须 status=completed 且 output_video_url 非空) cover_url: 可选的封面 URL 覆盖(finalize 时传入即优先使用,否则取 job.output_cover_url) """ from packages.adapters.sqlalchemy_impl.generated_video_repository import ( SQLAlchemyGeneratedVideoRepository, ) from packages.domain.generated_video import GeneratedVideo clip_name = f"AI数字人_{job.id[:8]}" # AI数字人入口是独立页面,前端可能不传 project_id(无项目概念), # 兜底为 "ai_avatar" 避免 DB 非空约束/查询问题;generation_task_id 用 render_job_id 便于反查。 clip_project_id = (job.project_id or "").strip() or "ai_avatar" clip_generation_task_id = job.id effective_cover = (cover_url or "").strip() if cover_url else (job.output_cover_url or "").strip() clip = GeneratedVideo.create( project_id=clip_project_id, generation_task_id=clip_generation_task_id, name=clip_name, file_url=job.output_video_url, user_id=job.user_id, duration=job.output_duration or 0.0, thumbnail_url=effective_cover or None, generation_params={ "source": "ai_avatar_render", "render_job_id": job.id, }, ) video_repo = SQLAlchemyGeneratedVideoRepository(self.db) video_repo.create(clip) logger.info("[数字人渲染] 成片已入库: clip_id=%s render_job=%s", clip.id, job.id) return clip def finalize_job(self, job_id: str, user_id: str, cover_url: Optional[str] = None): """用户在前端点「完成」后调用:将已 completed 的渲染任务正式入库到成片库. - 必须 status=completed 才可调用 - cover_url 若传入则优先使用并回写 job.output_cover_url;否则使用 job.output_cover_url(smart-cover/custom-cover 已写入) - 幂等:已入库则返回已存在的 GeneratedVideo """ from packages.adapters.sqlalchemy_impl.generated_video_repository import ( SQLAlchemyGeneratedVideoRepository, ) from packages.adapters.sqlalchemy_impl.models import GeneratedVideoModel job = self.get_render_job(job_id, user_id) if job is None: raise AiAvatarRenderError("渲染任务不存在", code="RenderJobNotFound") if job.status != "completed": raise AiAvatarRenderError(f"渲染任务未完成(当前状态: {job.status}),无法入库", code="RenderNotCompleted") if not (job.output_video_url or "").strip(): raise AiAvatarRenderError("渲染成片视频 URL 为空,无法入库", code="OutputVideoMissing") # 幂等检查:已入库直接返回现有记录(通过 generation_task_id=job_id 识别, # 因为入库时 generation_task_id 被设置为 render_job_id 自身) existing = ( self.db.query(GeneratedVideoModel) .filter( GeneratedVideoModel.user_id == user_id, GeneratedVideoModel.generation_task_id == job_id, ) .first() ) if existing is not None: logger.info("[数字人渲染] finalize 幂等命中,返回已存在记录: clip_id=%s job_id=%s", existing.id, job_id) return SQLAlchemyGeneratedVideoRepository(self.db).get(existing.id) # 传入 cover_url 时回写到 job if cover_url and cover_url.strip(): job.output_cover_url = cover_url.strip() # 同步更新 cover_config,保持 smart-cover 路径一致 if isinstance(job.cover_config, dict): job.cover_config = {**job.cover_config, "mode": "auto_frame", "url": cover_url.strip()} job.updated_at = datetime.now(UTC) self.db.commit() return self._persist_to_library(job, cover_url=cover_url) def _download_video(self, url: str) -> str: """下载视频到临时文件.""" import httpx tmp = tempfile.NamedTemporaryFile(suffix=".mp4", delete=False) try: with httpx.Client(timeout=120) as client: resp = client.get(url) resp.raise_for_status() tmp.write(resp.content) return tmp.name except Exception: if os.path.exists(tmp.name): os.unlink(tmp.name) raise @staticmethod def _save_title_dataurl_to_file(dataurl: str, *, dst_path: str | None = None, job_id: str = "") -> str: """解码前端传来的 data:image/png;base64,... 并保存为本地 PNG 文件。 Args: dataurl: 完整 dataURL 字符串 dst_path: 指定输出路径;为 None 时创建临时文件并返回路径 job_id: 仅在 dst_path 为空时用于临时文件命名 Returns: 保存后的本地文件路径 """ if not isinstance(dataurl, str) or not dataurl.startswith("data:image/"): raise ValueError("title_image_dataurl 不是合法的 data:image URL") # 拆分 data:image/png;base64, try: header, b64 = dataurl.split(",", 1) except ValueError as exc: raise ValueError("title_image_dataurl 缺少 base64 payload") from exc if "base64" not in header: raise ValueError("title_image_dataurl 不是 base64 编码") try: png_bytes = base64.b64decode(b64, validate=True) except (binascii.Error, ValueError) as exc: raise ValueError(f"title_image_dataurl base64 解码失败: {exc}") from exc if not png_bytes: raise ValueError("title_image_dataurl 解码后为空") if dst_path: out_path = dst_path with open(out_path, "wb") as f: f.write(png_bytes) return out_path suffix = f"_title_{job_id}.png" if job_id else "_title.png" with tempfile.NamedTemporaryFile(suffix=suffix, delete=False) as tmp: tmp.write(png_bytes) return tmp.name @staticmethod def _probe_video_resolution(video_path: str) -> tuple[int, int]: """用 ffprobe 探测视频分辨率,返回 (width, height);失败返回 (0, 0)。""" try: result = subprocess.run( [ "ffprobe", "-v", "error", "-select_streams", "v:0", "-show_entries", "stream=width,height", "-of", "csv=p=0:s=x", video_path, ], capture_output=True, text=True, timeout=15, ) if result.returncode == 0 and result.stdout.strip(): parts = result.stdout.strip().split("x") if len(parts) == 2: w, h = int(parts[0]), int(parts[1]) if w > 0 and h > 0: return w, h except Exception as exc: logger.warning("[数字人渲染] ffprobe 探测分辨率失败: %s", exc) return 0, 0 def _build_ffmpeg_command( self, *, input_video: str, b_roll_segments: list[dict[str, Any]], extra_inputs: list[str] | None = None, filter_complex: str, final_label: Optional[str], output_path: str, ) -> list[str]: """构建 FFmpeg 命令(list 形式,shell=False). 根因修复 #1798 P0:OSS 预签名 URL 含 `&Expires=...&Signature=...` 特殊字符, os.system(shell=True) 会把 `&` 解释为后台命令分隔符,导致 -filter_complex 被 当成独立命令报 sh: -filter_complex: not found(exit 127 → Python 32512)。 list + shell=False 彻底规避 shell 转义问题。 """ cmd: list[str] = ["ffmpeg", "-i", input_video] for seg in b_roll_segments: asset_url = seg.get("asset_url", "") if asset_url: cmd.extend(["-i", asset_url]) # 额外输入(例如前端 Canvas 渲染的标题 PNG) for extra in extra_inputs or []: cmd.extend(["-i", extra]) if filter_complex and final_label: cmd.extend( [ "-filter_complex", filter_complex, "-map", f"[{final_label}]", "-map", "0:a?", ] ) elif filter_complex: cmd.extend(["-filter_complex", filter_complex]) cmd.extend( [ "-c:v", "libx264", "-preset", "veryfast", "-crf", "23", "-c:a", "aac", "-b:a", "128k", "-y", output_path, ] ) return cmd def _upload_to_oss(self, local_path: str, oss_key: str) -> str: """上传文件到 OSS,返回 URL. 使用 SharedStorageService 统一存储服务。 """ storage = get_shared_storage_service() url = storage.upload_file_smart(local_path, oss_key) if url is None: raise AiAvatarRenderError( f"上传文件到 OSS 失败: {oss_key}", code="OSSUploadFailed", ) logger.info("上传文件到 OSS 成功: %s -> %s", local_path, url) return url