Files
xiaoxia-saas/apps/worker/worker_app/tasks/viral_video.py
T
saas-bot 8bdc39a1ab
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 2s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
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 / Frontend Lint (pull_request) Failing after 1m10s
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 / PR Build API Image (pull_request) Successful in 1m21s
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
CI/CD Pipeline / Frontend Unit Tests (pull_request) Successful in 1m43s
CI/CD Pipeline / PR Build Web Image (pull_request) Failing after 2m14s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m2s
Preview Deploy / Deploy Preview Environment (pull_request) Failing after 4m5s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 6m32s
AI Code Review / AI Code Review (pull_request) Successful in 7m14s
CI/CD Pipeline / Unit Tests (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 / Validate - Style (pull_request) Has been cancelled
CI/CD Pipeline / Validate - Security (pull_request) Has been cancelled
CI/CD Pipeline / Validate - Python (mypy + alembic) (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 / Integration Tests (pull_request) Has been cancelled
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been cancelled
fix: P0 upload 404 (project/library mismatch) + worker bug fixes
- upload.ts: library_id optional, auto ensureDefaultLibrary via getOrCreateDefaultProject
- 5 call sites simplified (ViralVideoPage/useVoiceUpload/useBatchCovers/Step6Cover/useVoiceMaterials)
- Bug1 (P0): upload_to_oss prefix kwarg → construct storage_key per generation.py pattern
- Bug2 (P0): Seedance 404 enhanced logging w/ URL/model/response body + troubleshooting hints
- Bug3 (P1): TTS import from services.* → apps.worker.services.* + synthesize fallback
- duplicated=true: backend returns existing asset URL in prepare, frontend uses it
- CI: ci_staging_deploy.sh auto-clean stale docker-compose.override.yml

Refs: #2109
2026-09-30 22:14:28 +08:00

842 lines
32 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.
"""爆款视频 Celery 编排器 — ViralVideoOrchestrator.
9 步流水线(Seedance 2.5 直生口型,不再走 MuseTalk):
1. 图片 VLM 分析
1.5 [v1.3] 视频风格分析(如用户上传参考视频)
2. 用户文案意图解析
3. 文案融合生成
4. 分镜脚本生成
5. 合规审核(6 维度,不通过自动重写 1 次)
6. CosyVoice 配音
7. BGM 选择(素材未就绪时跳过)
8. Seedance 逐分镜生成 + ffmpeg concat + 混 TTS
9. OSS 上传 + 通知 + 扣点
"""
from __future__ import annotations
import json
import logging
import os
from celery import Task, shared_task
from celery.exceptions import Retry
from worker_app.celery_app import celery_app # noqa: F401 - 加载 app 以注册任务
from worker_app.db import SessionLocal
from packages.adapters.sqlalchemy_impl.viral_video_repository import (
SQLAlchemyViralVideoJobRepository,
)
from packages.domain.viral_video import (
CREDITS_VIRAL_VIDEO_COST,
STAGE_LABELS,
ViralVideoJob,
ViralVideoStage,
ViralVideoStatus,
)
logger = logging.getLogger(__name__)
# ── WS 进度推送 ──────────────────────────────────────────────────────────
def _emit_progress(
job_id: str,
stage: str,
progress: float,
message: str = "",
data: dict | None = None,
event_type: str = "viral_video:progress",
):
"""通过 Redis 发布进度事件,供 WebSocket 消费。
event_type 取值:
- viral_video:progress 中间进度(默认)
- viral_video:completed 任务完成
- viral_video:failed 任务失败
- viral_video:wait_user 等待用户确认
所有事件 payload 均为合法 JSON,前端 JSON.parse 即可。
"""
try:
import redis as redis_lib
redis_url = os.environ.get("REDIS_URL", "redis://localhost:6379/0")
r = redis_lib.from_url(redis_url)
event = {
"type": event_type,
"job_id": job_id,
"stage": stage,
"progress": progress,
"message": message or STAGE_LABELS.get(stage, stage),
"data": data or {},
}
r.publish(f"viral_video:{job_id}", json.dumps(event, ensure_ascii=False))
except Exception as e:
logger.warning("[爆款视频] WS 进度推送失败: %s", e)
# ── 仓储辅助 ────────────────────────────────────────────────────────────
def _get_repo_and_job(job_id: str):
"""获取 session, repo, job 三元组。"""
session = SessionLocal()
repo = SQLAlchemyViralVideoJobRepository(session)
job = repo.get(job_id)
return session, repo, job
def _save_job(repo, job, session):
"""持久化并关闭 session。"""
repo.update(job)
session.commit()
# ── 流水线各步骤 ────────────────────────────────────────────────────────
# ── 流水线各步骤 ────────────────────────────────────────────────────────
def _step_image_analysis(job: ViralVideoJob) -> dict:
"""步骤 1: 图片 VLM 分析 — 识别产品特征、场景、卖点。"""
try:
from packages.shared.ai_service import call_vision
except ImportError:
logger.warning("[爆款视频] ai_service.call_vision 不可用,使用占位结果")
return {"products": [{"name": "产品", "features": ["特征1", "特征2"], "scene": "通用场景"}]}
results = []
for img_url in job.images:
try:
result = call_vision(
image_url=img_url,
prompt="请分析这张产品图片,识别:1)产品名称和类别 2)主要特征和卖点 3)适用场景 4)视觉风格。以JSON格式返回。",
)
results.append(result)
except Exception as e:
logger.warning("[爆款视频] 图片分析失败 img=%s: %s", img_url, e)
results.append({"name": "未识别", "features": [], "scene": "通用"})
return {"products": results}
def _step_video_analysis(job: ViralVideoJob) -> dict | None:
"""步骤 1.5 [v1.3]: 参考视频风格分析。"""
if not job.reference_video_url:
return None
try:
# P0-2: 修正 import 路径(video_analyzer.py 在 apps/worker/viral_video/ 下,worker PYTHONPATH 含 apps/worker)
from viral_video.video_analyzer import analyze_video_style
style_guide = analyze_video_style(job.reference_video_url)
return style_guide
except ImportError as e:
logger.info("[爆款视频] video_analyzer 模块未就绪(%s),使用占位风格分析", e)
return {
"cut_speed": "medium",
"transition": "cross_dissolve",
"energy": "medium",
"color_grade": "neutral",
"narrative": False,
"source": "placeholder",
}
except Exception as e:
logger.error("[爆款视频] 视频风格分析失败: %s", e)
return {"error": str(e), "source": "failed"}
def _step_intent_parsing(job: ViralVideoJob, image_analysis: dict) -> dict:
"""步骤 2: 用户文案意图解析 — 理解用户想表达什么。"""
try:
from packages.shared.ai_service import call_llm
except ImportError:
return {"intent": "推广产品", "key_messages": ["产品亮点"], "tone": "专业"}
products_summary = ""
for p in image_analysis.get("products", []):
products_summary += f"- {p.get('name', '产品')}: {', '.join(p.get('features', []))}\n"
prompt = f"""你是一个营销文案策略师。请分析以下信息,理解用户的营销意图:
用户原始文案:{job.user_copy_text or "(未提供)"}
行业:{job.industry or "未指定"}
目标客户:{job.target_customer or "未指定"}
营销目的:{job.marketing_purpose or "未指定"}
产品信息:
{products_summary}
请分析并返回JSON格式:
1. intent: 核心营销意图(一句话)
2. key_messages: 要传达的3-5个关键信息
3. tone: 文案调性(如:专业/亲切/高端/活力)
4. target_emotion: 希望触发的用户情感
5. call_to_action: 行动号召建议"""
try:
result = call_llm(prompt)
return result if isinstance(result, dict) else {"raw": result}
except Exception as e:
logger.warning("[爆款视频] 意图解析失败: %s", e)
return {"intent": "推广产品", "key_messages": ["产品亮点"], "tone": "专业"}
def _step_copy_fusion(job: ViralVideoJob, intent: dict, image_analysis: dict) -> str:
"""步骤 3: 文案融合生成 — 根据 fusion_level 融合用户文案和 AI 文案。"""
try:
from packages.shared.ai_service import call_llm
except ImportError:
return f"【{job.industry or '行业'}】优质产品,{job.target_customer or '您'}的不二之选!"
products_desc = ""
for p in image_analysis.get("products", []):
products_desc += f"{p.get('name', '产品')}({','.join(p.get('features', []))})\n"
if job.fusion_level == "ai_full":
prompt = f"""请为以下产品撰写一段爆款短视频文案({job.duration}秒):
产品:{products_desc}
行业:{job.industry}
目标客户:{job.target_customer}
营销目的:{job.marketing_purpose}
调性:{intent.get("tone", "专业")}
关键信息:{", ".join(intent.get("key_messages", []))}
要求:吸引眼球、节奏紧凑、有行动号召。直接输出文案内容。"""
elif job.fusion_level == "user_primary":
prompt = f"""请基于用户原始文案进行润色优化,保留用户原意和风格:
用户原文:{job.user_copy_text}
产品信息:{products_desc}
要求:保留用户原意,仅修正表达和节奏。直接输出文案内容。"""
else: # ai_polish (default)
prompt = f"""请将用户文案与AI分析融合,生成一段优化后的爆款短视频文案({job.duration}秒):
用户原文:{job.user_copy_text or "(未提供)"}
产品分析:{products_desc}
行业:{job.industry}
目标客户:{job.target_customer}
营销目的:{job.marketing_purpose}
意图分析:{intent.get("intent", "")}
调性:{intent.get("tone", "专业")}
要求:融合用户意图和产品卖点,节奏紧凑,适合短视频。直接输出文案内容。"""
try:
result = call_llm(prompt)
return result if isinstance(result, str) else str(result)
except Exception as e:
logger.warning("[爆款视频] 文案融合失败: %s", e)
return job.user_copy_text or f"精选{job.industry or '行业'}好物,值得关注!"
def _step_storyboard(job: ViralVideoJob, copy_text: str, image_analysis: dict) -> list[dict]:
"""步骤 4: 分镜脚本生成。每个分镜独立一段视频,段内时长建议 3~6 秒。"""
try:
from packages.shared.ai_service import call_llm
except ImportError:
return [
{
"order": 0,
"type": "product_shot",
"text": copy_text[:50],
"duration": min(5, job.duration),
"description": "产品展示",
"ken_burns": "zoom_in",
"transition": "cut",
}
]
products_hint = ""
products = image_analysis.get("products", []) if image_analysis else []
if products:
p0 = products[0] if isinstance(products[0], dict) else {}
feats = p0.get("features", []) if isinstance(p0, dict) else []
products_hint = f"\n首帧参考产品特征:{p0.get('name','')} - {', '.join(feats[:3])}"
seg_seconds = 5
n_segments = max(2, min(6, max(1, job.duration // seg_seconds)))
ratio = "9:16"
prompt = f"""请根据以下文案生成爆款短视频分镜脚本,共 {n_segments} 个分镜:
文案内容:{copy_text}
视频总时长:{job.duration}秒(每个分镜 3~6 秒,总和约等于总时长)
风格强度:{job.style_strength}
输出宽高比:{ratio}{products_hint}
请以 JSON 数组格式返回分镜列表,每个分镜包含:
- order: 序号(从0开始)
- type: 镜头类型(product_shot/close_up/scene/action/text_card/closing)
- description: 画面详细描述(中文,含主体、动作、场景、运镜、光影,用于AI视频生成prompt)
- text: 该分镜配音/字幕文本
- duration: 时长(秒,3~6秒的整数)
- ken_burns: 运镜方式(zoom_in/zoom_out/pan_left/pan_right/static)
- transition: 与下一分镜的转场(cut/dissolve/fade)"""
try:
result = call_llm(prompt)
if isinstance(result, list):
return _normalize_storyboard(result, job.duration, n_segments, copy_text)
import json as _json
parsed = _json.loads(result) if isinstance(result, str) else result
if isinstance(parsed, list):
return _normalize_storyboard(parsed, job.duration, n_segments, copy_text)
except Exception as e:
logger.warning("[爆款视频] 分镜生成失败: %s", e)
return _fallback_storyboard(copy_text, job.duration, n_segments)
def _normalize_storyboard(raw: list, total_duration: int, n_segments: int, copy_text: str) -> list[dict]:
"""规范化 LLM 输出的分镜:填充缺省字段、保证总时长合理。"""
out: list[dict] = []
for i, item in enumerate(raw):
if not isinstance(item, dict):
continue
try:
dur = int(item.get("duration") or 5)
except (TypeError, ValueError):
dur = 5
dur = max(3, min(8, dur))
out.append(
{
"order": int(item.get("order", i)),
"type": str(item.get("type", "product_shot")),
"description": str(item.get("description", copy_text[:80])),
"text": str(item.get("text", "")),
"duration": dur,
"ken_burns": str(item.get("ken_burns", "zoom_in")),
"transition": str(item.get("transition", "cut")),
}
)
if not out:
return _fallback_storyboard(copy_text, total_duration, n_segments)
out = out[:n_segments]
total = sum(s["duration"] for s in out)
if total > 0 and total != total_duration:
scale = total_duration / total
acc = 0
for s in out[:-1]:
s["duration"] = max(3, min(8, round(s["duration"] * scale)))
acc += s["duration"]
out[-1]["duration"] = max(3, total_duration - acc)
return out
def _fallback_storyboard(copy_text: str, total_duration: int, n_segments: int) -> list[dict]:
if n_segments <= 0:
n_segments = 1
dur = total_duration // n_segments
remainder = total_duration - dur * n_segments
out = []
for i in range(n_segments):
d = dur + (remainder if i == n_segments - 1 else 0)
out.append(
{
"order": i,
"type": "product_shot",
"description": f"产品展示镜头 {i + 1}:{copy_text[:40]}",
"text": copy_text,
"duration": max(3, d),
"ken_burns": "zoom_in" if i % 2 == 0 else "pan_left",
"transition": "cut",
}
)
return out
def _step_review(job: ViralVideoJob, copy_text: str, storyboard: list[dict]) -> dict:
"""步骤 5: 合规审核(6 维度)。不通过时自动重写 1 次。"""
dimensions = ["广告法合规", "平台规范", "内容真实性", "版权安全", "价值观", "风格一致性"]
try:
from packages.shared.ai_service import call_llm
except ImportError:
return {"passed": True, "score": 90, "details": {d: "通过" for d in dimensions}}
prompt = f"""请对以下短视频内容进行合规审核,检查6个维度:{", ".join(dimensions)}
文案内容:{copy_text}
分镜脚本:{storyboard[:3]}...
行业:{job.industry}
请以JSON格式返回:
- passed: bool(是否全部通过)
- score: int(0-100分)
- details: 各维度评分和说明
- issues: 需要修改的问题列表(如有)"""
try:
result = call_llm(prompt)
return result if isinstance(result, dict) else {"passed": True, "score": 80, "details": {}}
except Exception as e:
logger.warning("[爆款视频] 合规审核失败: %s", e)
return {"passed": True, "score": 75, "details": {d: "默认通过" for d in dimensions}}
def _step_tts(job: ViralVideoJob, copy_text: str):
"""步骤 6: CosyVoice 配音。P1:返回 Path;失败返回 None。"""
try:
from pathlib import Path as _Path
# 使用绝对包路径,避免 celery worker 因 cwd/PYTHONPATH 微小差异找不到 services 模块
from apps.worker.services.tts_service_factory import get_tts_service
tts_service = get_tts_service()
# 兼容老接口:部分 provider 只接收 text 参数
try:
result = tts_service.synthesize(text=copy_text, voice_id=job.persona_id or "default")
except TypeError:
result = tts_service.synthesize(text=copy_text)
if result is None:
return None
p = _Path(result) if not isinstance(result, _Path) else result
if p.exists():
return p
logger.warning("[爆款视频] TTS 返回路径不存在: %s", p)
return None
except Exception as e:
logger.warning("[爆款视频] TTS 配音失败: %s", e)
return None
def _step_bgm_select(job: ViralVideoJob):
"""步骤 7: BGM 选择。P1:素材未就绪前返回 None,跳过 BGM 混音。"""
return None
def _build_segment_prompt(seg: dict, job: ViralVideoJob, style_hint: str) -> str:
desc = seg.get("description") or seg.get("text") or "产品展示"
ken_burns = seg.get("ken_burns", "zoom_in")
cam_map = {
"zoom_in": "缓慢推镜放大",
"zoom_out": "缓慢拉镜缩小",
"pan_left": "镜头向左平移",
"pan_right": "镜头向右平移",
"static": "固定镜头",
}
camera = cam_map.get(ken_burns, "缓慢运镜")
parts = [
f"{desc}。",
f"运镜:{camera}。",
"画面流畅、电影感光影、高清细节,9:16竖屏,适合短视频。",
]
if style_hint:
parts.append(f"参考风格:{style_hint}")
return " ".join(parts)
def _step_render(job, storyboard, tts_path, bgm):
"""步骤 8: 渲染(P0-1 核心重写)。
每个 storyboard 分镜 → Seedance 2.5 生成短视频段(无声)→ 下载 → ffmpeg concat → 混入 TTS。
返回最终视频本地路径字符串。
"""
import tempfile
from pathlib import Path
from video_processing.concat_engine import concat_video_files
from packages.shared.ai_service import call_video_generation
from packages.shared.ffmpeg_utils import run_ffmpeg
if not storyboard:
raise ValueError("storyboard is empty")
style_hint = ""
if isinstance(job.style_guide, dict):
style_hint = f"节奏{job.style_guide.get('cut_speed','')}、转场{job.style_guide.get('transition','')}、色调{job.style_guide.get('color_grade','')}"
tmpdir = Path(tempfile.mkdtemp(prefix=f"viral_{job.id}_"))
logger.info("[爆款视频] 开始渲染,分镜数=%d, tmpdir=%s", len(storyboard), tmpdir)
seg_paths: list[str] = []
first_image = job.images[0] if job.images else None
n_total = len(storyboard)
for i, seg in enumerate(storyboard):
try:
dur = int(seg.get("duration") or 5)
except (TypeError, ValueError):
dur = 5
dur = max(2, min(12, dur))
prompt = _build_segment_prompt(seg, job, style_hint)
_emit_progress(
job.id,
ViralVideoStage.RENDERING,
80.0 + (i + 1) / max(n_total, 1) * 5.0,
f"正在生成分镜 {i + 1}/{n_total} ({dur}s)...",
)
logger.info("[爆款视频] 分镜 %d/%d dur=%ds prompt=%s", i + 1, n_total, dur, prompt[:80])
seg_path = call_video_generation(
prompt=prompt,
image_url=first_image if i == 0 else None,
duration=dur,
ratio="9:16",
resolution="720p",
output_dir=str(tmpdir),
)
if not seg_path or not Path(seg_path).exists():
logger.warning("[爆款视频] 分镜 %d 生成失败,使用占位片段", i + 1)
seg_path = str(_make_placeholder_clip(tmpdir, i, dur))
seg_paths.append(seg_path)
_emit_progress(job.id, ViralVideoStage.RENDERING, 86.0, "正在拼接分镜...")
concat_out = tmpdir / "concat_raw.mp4"
try:
concat_video_files(seg_paths, concat_out, work_dir=tmpdir, force_reencode=True)
except Exception as e:
logger.error("[爆款视频] concat 失败: %s,降级过滤无效片段", e, exc_info=True)
valid = [p for p in seg_paths if _probe_ok(p)]
if not valid:
raise RuntimeError(f"所有分镜片段均无效: {e}") from e
concat_video_files(valid, concat_out, work_dir=tmpdir, force_reencode=True)
final_path = concat_out
if tts_path is not None:
tts_p = Path(tts_path) if not isinstance(tts_path, Path) else tts_path
if tts_p.exists():
_emit_progress(job.id, ViralVideoStage.RENDERING, 87.5, "正在合成配音...")
mixed_out = tmpdir / "final_with_audio.mp4"
try:
run_ffmpeg(
[
"ffmpeg",
"-y",
"-i",
str(concat_out),
"-i",
str(tts_p),
"-c:v",
"copy",
"-c:a",
"aac",
"-b:a",
"192k",
"-map",
"0:v:0",
"-map",
"1:a:0",
"-shortest",
str(mixed_out),
]
)
if mixed_out.exists() and mixed_out.stat().st_size > 0:
final_path = mixed_out
except Exception as e:
logger.warning("[爆款视频] TTS 混音失败,使用无声视频: %s", e)
logger.info("[爆款视频] 渲染完成: %s size=%d", final_path, final_path.stat().st_size if final_path.exists() else 0)
return str(final_path)
def _probe_ok(video_path: str) -> bool:
import subprocess
from pathlib import Path as _Path
try:
if not _Path(video_path).exists():
return False
r = subprocess.run(
[
"ffprobe",
"-v",
"error",
"-select_streams",
"v:0",
"-show_entries",
"stream=codec_type",
"-of",
"csv=p=0",
video_path,
],
capture_output=True,
timeout=10,
)
return r.returncode == 0 and b"video" in r.stdout
except Exception:
return False
def _make_placeholder_clip(tmpdir, idx: int, duration: int):
import subprocess
out = tmpdir / f"placeholder_{idx}.mp4"
try:
subprocess.run(
[
"ffmpeg",
"-y",
"-f",
"lavfi",
"-i",
f"color=c=0x202030:s=720x1280:d={max(duration,2)}:r=24",
"-f",
"lavfi",
"-i",
f"anullsrc=r=44100:cl=stereo:d={max(duration,2)}",
"-c:v",
"libx264",
"-pix_fmt",
"yuv420p",
"-preset",
"ultrafast",
"-c:a",
"aac",
"-shortest",
str(out),
],
capture_output=True,
timeout=60,
check=True,
)
except Exception as e:
logger.warning("[爆款视频] 占位片段生成失败: %s", e)
return out
def _step_upload(job: ViralVideoJob, video_path: str) -> str:
"""步骤 9: OSS 上传。"""
from pathlib import Path
from video_processing.oss_helpers import upload_to_oss
local = Path(video_path)
# 构造 OSS key,与 generation.py 规则对齐:generated/viral-video/<user_id>/<job_id>/<filename>
storage_key = f"generated/viral-video/{job.user_id}/{job.id}/{local.name}"
logger.info("[爆款视频] 开始上传成片: local=%s key=%s size=%d", local, storage_key, local.stat().st_size)
video_url = upload_to_oss(local, storage_key)
if not video_url:
raise RuntimeError(f"OSS 上传失败: storage_key={storage_key}")
return video_url
# ── 主编排器 ────────────────────────────────────────────────────────────
@shared_task(bind=True, max_retries=2, name="worker.run_viral_video_pipeline")
def run_viral_video_pipeline(self: Task, job_id: str) -> dict:
"""爆款视频 10 步流水线编排器(前半段:图片分析→风格分析→意图解析,然后 WAIT_USER_CONFIRM)。"""
session = None
try:
session, repo, job = _get_repo_and_job(job_id)
if job is None:
logger.error("[爆款视频] 任务不存在: %s", job_id)
return {"ok": False, "error": "job not found"}
job.mark_running()
_save_job(repo, job, session)
_emit_progress(job_id, ViralVideoStage.IMAGE_ANALYSIS, 5.0, "开始图片分析")
# ── Step 1: 图片 VLM 分析 ──
_emit_progress(job_id, ViralVideoStage.IMAGE_ANALYSIS, 10.0, "正在分析产品图片...")
image_analysis = _step_image_analysis(job)
# P0-3: 持久化 image_analysis 到 job,供 resume 阶段使用
job.image_analysis = image_analysis
_save_job(repo, job, session)
_emit_progress(job_id, ViralVideoStage.IMAGE_ANALYSIS, 15.0, "图片分析完成", {"result": image_analysis})
# ── Step 1.5: 视频风格分析(v1.3) ──
style_guide = None
if job.reference_video_url or job.style_template_id:
_emit_progress(job_id, ViralVideoStage.VIDEO_ANALYSIS, 20.0, "正在分析参考视频风格...")
style_guide = _step_video_analysis(job)
job.style_guide = style_guide
_save_job(repo, job, session)
_emit_progress(
job_id,
ViralVideoStage.VIDEO_ANALYSIS,
25.0,
"风格分析完成",
{"style_analyzed": True, "style_guide": style_guide},
)
# ── Step 2: 意图解析 ──
_emit_progress(job_id, ViralVideoStage.INTENT_PARSING, 30.0, "正在解析文案意图...")
intent_result = _step_intent_parsing(job, image_analysis)
job.mark_wait_user_confirm(intent_result)
_save_job(repo, job, session)
_emit_progress(
job_id,
ViralVideoStage.INTENT_PARSING,
35.0,
"意图解析完成,等待用户确认",
{"intent_result": intent_result, "waiting_confirm": True},
)
_emit_progress(
job_id,
ViralVideoStage.INTENT_PARSING,
35.0,
"等待用户确认意图文案",
{"intent_result": intent_result},
event_type="viral_video:wait_user",
)
return {"ok": True, "job_id": job_id, "status": "wait_user_confirm", "intent_result": intent_result}
except Retry:
raise
except Exception as e:
logger.error("[爆款视频] 流水线异常: %s", e, exc_info=True)
err_msg = str(e)
failed_stage = ""
try:
if session is None:
session = SessionLocal()
repo = SQLAlchemyViralVideoJobRepository(session)
job = repo.get(job_id)
else:
_, repo, job = _get_repo_and_job(job_id)
if job is not None and not job.is_terminal:
job.mark_failed(err_msg)
failed_stage = getattr(job, "current_stage", "") or ""
_save_job(repo, job, session)
except Exception as inner:
logger.warning("[爆款视频] 标记失败状态时出错: %s", inner)
_emit_progress(
job_id,
failed_stage,
0,
f"任务失败: {err_msg}",
{"error": err_msg},
event_type="viral_video:failed",
)
return {"ok": False, "job_id": job_id, "error": err_msg}
finally:
if session:
session.close()
@shared_task(bind=True, max_retries=2, name="worker.resume_viral_video_pipeline")
def resume_viral_video_pipeline(self: Task, job_id: str) -> dict:
"""用户确认意图后,从断点恢复流水线(步骤 3-10)。"""
session = None
job = None
try:
session, repo, job = _get_repo_and_job(job_id)
if job is None:
return {"ok": False, "error": "job not found"}
if job.status != ViralVideoStatus.RUNNING:
return {"ok": False, "error": f"unexpected status: {job.status}"}
# P0-3: 从 job 读取 image_analysis(run_pipeline 阶段已持久化)
image_analysis = job.image_analysis or {"products": []}
_emit_progress(job_id, ViralVideoStage.COPY_FUSION, 40.0, "正在融合文案...")
# ── Step 3: 文案融合 ──
copy_text = _step_copy_fusion(job, job.intent_result or {}, image_analysis)
_emit_progress(job_id, ViralVideoStage.COPY_FUSION, 50.0, "文案融合完成")
# ── Step 4: 分镜脚本 ──
_emit_progress(job_id, ViralVideoStage.STORYBOARD, 55.0, "正在生成分镜脚本...")
storyboard = _step_storyboard(job, copy_text, image_analysis)
_emit_progress(job_id, ViralVideoStage.STORYBOARD, 60.0, "分镜脚本完成", {"segments": len(storyboard)})
# ── Step 5: 合规审核 ──
_emit_progress(job_id, ViralVideoStage.REVIEW, 65.0, "正在进行合规审核...")
review_result = _step_review(job, copy_text, storyboard)
if not review_result.get("passed", True):
_emit_progress(job_id, ViralVideoStage.REVIEW, 67.0, "审核未通过,正在自动重写...")
copy_text = _step_copy_fusion(job, job.intent_result or {}, image_analysis)
review_result = _step_review(job, copy_text, storyboard)
_emit_progress(job_id, ViralVideoStage.REVIEW, 70.0, "合规审核完成")
# ── Step 6: CosyVoice 配音(返回 Path | None) ──
_emit_progress(job_id, ViralVideoStage.TTS, 72.0, "正在生成配音...")
tts_path = _step_tts(job, copy_text)
_emit_progress(job_id, ViralVideoStage.TTS, 75.0, "配音完成", {"has_tts": tts_path is not None})
# ── Step 7: BGM 选择(P1:暂返回 None,跳过) ──
_emit_progress(job_id, ViralVideoStage.BGM_SELECT, 77.0, "BGM 已跳过(素材未就绪)")
bgm = _step_bgm_select(job)
# ── Step 8: 渲染(逐分镜 Seedance → concat → 混 TTS) ──
_emit_progress(job_id, ViralVideoStage.RENDERING, 80.0, "正在渲染视频...")
video_path = _step_render(job, storyboard, tts_path, bgm)
_emit_progress(job_id, ViralVideoStage.RENDERING, 88.0, "渲染完成")
# ── Step 9: OSS 上传 + 扣点 ──
# 注:爆款视频由 Seedance 2.5 直接生成口型,不需要 MuseTalk 事后对口型(MuseTalk 是 AI 数字人路线用的)。
_emit_progress(job_id, ViralVideoStage.UPLOADING, 95.0, "正在上传视频...")
video_url = _step_upload(job, video_path)
job.credits_cost = CREDITS_VIRAL_VIDEO_COST
# TODO: 调用 credits.deduct() 实际扣点(#1895 总开关为 false 时不扣,保留 TODO)
job.mark_completed(video_url)
_save_job(repo, job, session)
_emit_progress(job_id, ViralVideoStage.UPLOADING, 100.0, "视频生成完成!", {"video_url": video_url})
_emit_progress(
job_id,
ViralVideoStage.UPLOADING,
100.0,
"视频生成完成",
{"video_url": video_url},
event_type="viral_video:completed",
)
logger.info("[爆款视频] 任务完成: job_id=%s video_url=%s", job_id, video_url)
return {"ok": True, "job_id": job_id, "video_url": video_url}
except Retry:
raise
except Exception as e:
logger.error("[爆款视频] 恢复流水线异常: %s", e, exc_info=True)
err_msg = str(e)
failed_stage = ""
try:
if session is None:
session = SessionLocal()
repo = SQLAlchemyViralVideoJobRepository(session)
job = repo.get(job_id)
elif job is not None and not job.is_terminal:
job.mark_failed(err_msg)
failed_stage = getattr(job, "current_stage", "") or ""
_save_job(repo, job, session)
except Exception as inner:
logger.warning("[爆款视频] 标记失败状态时出错: %s", inner)
_emit_progress(
job_id,
failed_stage,
0,
f"任务失败: {err_msg}",
{"error": err_msg},
event_type="viral_video:failed",
)
return {"ok": False, "job_id": job_id, "error": err_msg}
finally:
if session:
session.close()
@shared_task(bind=True, max_retries=1, name="worker.run_video_style_analysis")
def run_video_style_analysis(self: Task, job_id: str) -> dict:
"""独立的视频风格分析任务(v1.3)。"""
session = None
try:
session, repo, job = _get_repo_and_job(job_id)
if job is None:
return {"ok": False, "error": "job not found"}
_emit_progress(job_id, ViralVideoStage.VIDEO_ANALYSIS, 10.0, "正在分析参考视频风格...")
style_guide = _step_video_analysis(job)
job.style_guide = style_guide
_save_job(repo, job, session)
_emit_progress(job_id, ViralVideoStage.VIDEO_ANALYSIS, 100.0, "风格分析完成", {"style_guide": style_guide})
return {"ok": True, "job_id": job_id, "style_guide": style_guide}
except Retry:
raise
except Exception as e:
logger.error("[爆款视频] 风格分析失败: %s", e)
return {"ok": False, "job_id": job_id, "error": str(e)}
finally:
if session:
session.close()