From fb10393f3df1881d091e90475e19e1b77c8ba92e Mon Sep 17 00:00:00 2001 From: xiaoxia-saas-bot Date: Sun, 4 Oct 2026 19:24:00 +0800 Subject: [PATCH] =?UTF-8?q?fix(viral=5Fvideo):=20#2040=20rebase=20?= =?UTF-8?q?=E5=90=8E=E6=B8=85=E7=90=86F841=E6=AD=BB=E4=BB=A3=E7=A0=81+?= =?UTF-8?q?=E4=BF=9D=E7=95=99#2180=20timeout=E4=BF=AE=E5=A4=8D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 删除 script_generation 中旧prompt遗留的未使用变量(style_hint/ratio/viral_structure_block/persona_hint) - 冲突解决保留 #2180 全部修复: VLM/LLM timeout 15/25→45, pro 25→60, max_retries 2→1 - black/isort/ruff 全绿;viral 相关71测试全过 --- apps/worker/worker_app/tasks/viral_video.py | 193 +++++++++++++++----- 1 file changed, 145 insertions(+), 48 deletions(-) diff --git a/apps/worker/worker_app/tasks/viral_video.py b/apps/worker/worker_app/tasks/viral_video.py index 4e7c572b7..54dbc16ae 100644 --- a/apps/worker/worker_app/tasks/viral_video.py +++ b/apps/worker/worker_app/tasks/viral_video.py @@ -50,6 +50,7 @@ logger = logging.getLogger(__name__) # ── WS 进度推送 ────────────────────────────────────────────────────────── + def _emit_progress( job_id: str, stage: str, @@ -75,18 +76,22 @@ def _emit_progress( except Exception as e: logger.warning("[爆款视频] WS 进度推送失败: %s", e) + # ── 仓储辅助 ──────────────────────────────────────────────────────────── + def _get_repo_and_job(job_id: str): session = SessionLocal() repo = SQLAlchemyViralVideoJobRepository(session) job = repo.get(job_id) return session, repo, job + def _save_job(repo, job, session): repo.update(job) session.commit() + def _start_trust_chain_preheat(job_id: str, portrait_descriptions: list[str]) -> None: """#2172/#2174 后台启动信任链预热(Seedream t2i 文生图人像),不阻塞调用方。 @@ -152,6 +157,7 @@ def _start_trust_chain_preheat(job_id: str, portrait_descriptions: list[str]) -> t = threading.Thread(target=_run_preheat, name=f"tc-preheat-{job_id[:8]}", daemon=True) t.start() + def _set_stage(job, repo, session, stage: str, message: str, persist: bool = True) -> None: """更新细粒度阶段并持久化到 DB,同时通过 Redis 推送进度事件。 @@ -167,6 +173,7 @@ def _set_stage(job, repo, session, stage: str, message: str, persist: bool = Tru except Exception as e: # 阶段持久化失败不阻塞主流程 logger.warning("[爆款视频] 阶段持久化失败 stage=%s err=%s", stage, e) + # ── worker 心跳(僵尸任务检测) ───────────────────────────────────────── # 心跳间隔(秒);超过此时间未更新 heartbeat_at 视为 worker 异常 @@ -176,6 +183,7 @@ _STALE_RUNNING_TIMEOUT_SEC = 10 * 60 # 10 分钟 # 心跳过期窗口:heartbeat_at 距 now 超过此时长视为失效 _HEARTBEAT_EXPIRE_SEC = 2 * 60 # 2 分钟 + def _heartbeat_once(job_id: str) -> None: """在独立 session 中更新一次 heartbeat_at(不捕获主流程事务状态)。""" ssn = None @@ -200,6 +208,7 @@ def _heartbeat_once(job_id: str) -> None: except Exception: pass + def _start_heartbeat_thread(job_id: str) -> tuple[threading.Event, threading.Thread]: """启动后台心跳线程,每 _HEARTBEAT_INTERVAL_SEC 秒更新一次 heartbeat_at。 返回 (stop_event, thread);任务结束时调用 stop_event.set() 停止心跳。 @@ -216,6 +225,7 @@ def _start_heartbeat_thread(job_id: str) -> tuple[threading.Event, threading.Thr t.start() return stop, t + def _recover_stale_jobs() -> int: """启动/定时扫描:把僵尸任务(running 超时且心跳停止)标记为 failed。 返回本次回收的任务数。可由 celery beat 周期性调用,也可在任务启动前顺带扫一次。 @@ -253,6 +263,7 @@ def _recover_stale_jobs() -> int: except Exception: pass + # ── 默认结构 ───────────────────────────────────────────────────────────── _DEFAULT_HARD_CONSTRAINTS = [ @@ -283,6 +294,7 @@ _DEFAULT_NEGATIVE_PROMPTS = [ "低分辨率", ] + def _empty_copy_result(duration: int = 15, ratio: str = "9:16") -> dict: return { "overview": {"theme": "好物推荐", "total_duration": duration, "aspect_ratio": ratio}, @@ -296,8 +308,10 @@ def _empty_copy_result(duration: int = 15, ratio: str = "9:16") -> dict: "title": "", } + # ── 流水线各步骤 ──────────────────────────────────────────────────────── + def _vision_fallback(idx: int, reason: str, extra: dict | None = None) -> dict: d = { "name": "未识别", @@ -315,6 +329,7 @@ def _vision_fallback(idx: int, reason: str, extra: dict | None = None) -> dict: d.update(extra) return d + def _is_vision_result_usable(result: dict) -> bool: """判断 VLM 返回是否有效:name/summary 不能为未识别/无法判断/空,summary 要够长。""" if not isinstance(result, dict): @@ -333,6 +348,7 @@ def _is_vision_result_usable(result: dict) -> bool: return False return True + def _analyze_single_image( idx: int, img_url: str, @@ -346,11 +362,13 @@ def _analyze_single_image( lite 失败/不可用时用 pro 降级重试 1 次。失败/None 最终返回含默认字段的 dict。 """ try: - from packages.shared.ai_service import call_vision - from packages.application.viral_video.prompt_loader import ( - get_template, render_system_prompt, render_user_prompt, - ) from packages.application.viral_video import xml_parser as xp + from packages.application.viral_video.prompt_loader import ( + get_template, + render_system_prompt, + render_user_prompt, + ) + from packages.shared.ai_service import call_vision except ImportError as e: logger.warning("[爆款视频] prompt 模板/解析模块不可用: %s", e) return _vision_fallback(idx, f"fallback_import_error:{e}") @@ -361,15 +379,21 @@ def _analyze_single_image( template = get_template("image_analysis") system = render_system_prompt(template) user = render_user_prompt( - template, image_count=1, industry="通用", + template, + image_count=1, + industry="通用", image_urls=f"第1张:{img_url}", ) def _call(model: str, tmo: int): try: return call_vision( - image_url=img_url, prompt=user, model=model, - max_tokens=2048, temperature=0.3, timeout=tmo, + image_url=img_url, + prompt=user, + model=model, + max_tokens=2048, + temperature=0.3, + timeout=tmo, system_prompt=system, ) except Exception as e: @@ -396,13 +420,18 @@ def _analyze_single_image( packaging = a.get("packaging", "") or "无法判断" summary = a.get("summary", "") or f"{brand} {name}" return { - "name": name, "brand": brand, "category": category, - "appearance": appearance, "packaging": packaging, + "name": name, + "brand": brand, + "category": category, + "appearance": appearance, + "packaging": packaging, "text_on_package": text_list, "key_features": feat_list or [features] if features else ["无法判断"], - "scene": scene, "mood": mood, + "scene": scene, + "mood": mood, "portrait_prompt": a.get("portrait_prompt", "无人像"), - "summary": summary, "_source": "xml", + "summary": summary, + "_source": "xml", } return _vision_fallback(idx, "no_product_tag") @@ -435,6 +464,7 @@ def _analyze_single_image( return pro_result return first_result + def _step_image_analysis(job: ViralVideoJob) -> dict: """步骤 1: 图片 VLM 分析 — 识别产品特征(v1.6 优化:并行 + lite 模型提速)。""" try: @@ -490,6 +520,7 @@ def _step_image_analysis(job: ViralVideoJob) -> dict: return {"products": results} + def _step_video_analysis(job: ViralVideoJob) -> dict | None: """步骤 1.5: 参考视频风格分析(可选)。""" if not job.reference_video_url: @@ -513,14 +544,17 @@ def _step_video_analysis(job: ViralVideoJob) -> dict | None: logger.error("[爆款视频] 视频风格分析失败: %s", e) return {"error": str(e), "source": "failed"} + def _step_intent_parsing(job: ViralVideoJob, image_analysis: dict) -> dict: """步骤 2: 用户文案意图解析(#2040:改为模板 + XML 解析)。""" try: - from packages.shared.ai_service import call_llm - from packages.application.viral_video.prompt_loader import ( - get_template, render_system_prompt, render_user_prompt, - ) from packages.application.viral_video import xml_parser as xp + from packages.application.viral_video.prompt_loader import ( + get_template, + render_system_prompt, + render_user_prompt, + ) + from packages.shared.ai_service import call_llm except ImportError: return {"intent": "推广产品", "key_messages": ["产品亮点"], "tone": "专业", "suggested_title": ""} @@ -556,11 +590,21 @@ def _step_intent_parsing(job: ViralVideoJob, image_analysis: dict) -> dict: msgs = [n["text"] for n in xp.find_all(raw, "message") if n["text"]] tone = xp.text_of(raw, "emotion_tone") or "亲切自然" title = xp.text_of(raw, "suggested_title") or xp.text_of(raw, "title") - return {"intent": summary or "推广产品", "key_messages": msgs or ["产品亮点"], "tone": tone, "suggested_title": title} + return { + "intent": summary or "推广产品", + "key_messages": msgs or ["产品亮点"], + "tone": tone, + "suggested_title": title, + } def _fallback(raw: str) -> dict: t = (job.user_copy_text or "").strip() - return {"intent": t[:30] or "推广产品", "key_messages": [t[:80]] if t else ["产品亮点"], "tone": "专业", "suggested_title": ""} + return { + "intent": t[:30] or "推广产品", + "key_messages": [t[:80]] if t else ["产品亮点"], + "tone": "专业", + "suggested_title": "", + } _s = get_shared_settings() _fast = _s.doubao_fast_model @@ -568,8 +612,13 @@ def _step_intent_parsing(job: ViralVideoJob, image_analysis: dict) -> dict: for _m, _lbl in [(_fast, "fast"), (_pro, "pro-fallback")]: try: logger.info("[爆款视频] 意图解析 model=%s label=%s", _m, _lbl) - raw = call_llm([{"role": "system", "content": system}, {"role": "user", "content": user}], - temperature=0.4, max_tokens=1024, model=_m, timeout=45) # #2180: 意图解析 LLM 实测需更长响应,原25s太紧 + raw = call_llm( + [{"role": "system", "content": system}, {"role": "user", "content": user}], + temperature=0.4, + max_tokens=1024, + model=_m, + timeout=45, + ) # #2180: 意图解析 LLM 实测需更长响应,原25s太紧 if not raw: continue parsed = _parse(raw) @@ -579,6 +628,7 @@ def _step_intent_parsing(job: ViralVideoJob, image_analysis: dict) -> dict: logger.warning("[爆款视频] 意图解析失败 label=%s err=%s", _lbl, e) return _fallback("") + _PERSONA_STYLE_GUIDE = { "通用个人IP": "亲切自然、像朋友分享好物,第一人称口语化,不端着", "老板型IP": "沉稳大气、有行业格局感,适度使用『我做了XX年』『我一直坚持』等老板视角,语气自信不夸张", @@ -592,6 +642,7 @@ _PERSONA_STYLE_GUIDE = { "测评种草型": "真实测评感、讲使用体验和优缺点对比,带『亲测』『我用了XX天』『实测下来』真实感词汇", } + def _persona_style_hint(persona_id: str) -> str: """根据 persona_id 查文案风格指导;未命中/空值返回通用提示。""" pid = (persona_id or "").strip() @@ -602,6 +653,7 @@ def _persona_style_hint(persona_id: str) -> str: return f"【人设风格:{pid}】按该人设的口吻、话术习惯组织口播和出镜动作" return "【人设风格:未指定】亲切自然、像朋友分享好物" + def _build_products_summary(image_analysis: dict) -> str: """把 VLM 返回的商品分析结果拼给文案/分镜生成 prompt 用。 优先用 summary(自然段落);没有时用结构化字段兜底拼一段。""" @@ -670,6 +722,7 @@ def _build_products_summary(image_analysis: dict) -> str: lines.append("- " + ",".join(parts)) return "\n".join(lines) + def _safe_json_loads(raw: str | dict | list | None): if raw is None: return None @@ -696,6 +749,7 @@ def _safe_json_loads(raw: str | dict | list | None): pass return None + def _replace_henjin_everywhere(obj: Any) -> Any: """递归遍历 copy_result 里所有字符串值,把'很近'替换成'最近'。 覆盖 overview.theme、scene_and_lighting、voiceover_script、 @@ -711,6 +765,7 @@ def _replace_henjin_everywhere(obj: Any) -> Any: return {k: _replace_henjin_everywhere(v) for k, v in obj.items()} return obj + def _fallback_script(job: ViralVideoJob) -> dict: """脚本生成失败时的兜底脚本(极简但可用)。""" dur = max(5, min(30, int(getattr(job, "duration", 15) or 15))) @@ -735,6 +790,7 @@ def _fallback_script(job: ViralVideoJob) -> dict: base["title"] = "好物分享" return base + def _validate_and_normalize_script(raw, job: ViralVideoJob) -> dict: """把 LLM 返回的脚本规范化、补默认、校验结构。""" dur = max(5, min(30, int(getattr(job, "duration", 15) or 15))) @@ -831,9 +887,11 @@ def _validate_and_normalize_script(raw, job: ViralVideoJob) -> dict: base["title"] = base["overview"]["theme"] return base + def _script_from_xml(raw: str, job: ViralVideoJob) -> dict | None: """把 LLM 返回的 XML 分镜规范化为旧 copy_result 结构(供 Seedance 使用)。""" from packages.application.viral_video import xml_parser as xp + dur = max(5, min(30, int(getattr(job, "duration", 15) or 15))) ratio = getattr(job, "video_ratio", None) or "9:16" base = _empty_copy_result(dur, ratio) @@ -860,14 +918,18 @@ def _script_from_xml(raw: str, job: ViralVideoJob) -> dict | None: ref_idx = xp.attr_int(ref_idx_raw, 0) shot = { "time_range": a.get("time_range") or f"{i * 3}-{(i + 1) * 3}秒", - "shot_type_angle_movement": (xp.text_of(body, "shot_type_angle_movement") if body else "") or a.get("shot_type_angle_movement", "") or "中景平视,固定镜头", + "shot_type_angle_movement": (xp.text_of(body, "shot_type_angle_movement") if body else "") + or a.get("shot_type_angle_movement", "") + or "中景平视,固定镜头", "scene_and_dialogue": (xp.text_of(body, "scene_and_dialogue") if body else "") or "", "action_details": (xp.text_of(body, "action_details") if body else "") or "", "audio_bgm": (xp.text_of(body, "audio_bgm") if body else "") or a.get("bgm_note", "") or "轻快BGM", - "transition": (xp.text_of(body, "transition") if body else "") or a.get("transition", "") or ("硬切" if i < len(clips) - 1 else "结束"), + "transition": (xp.text_of(body, "transition") if body else "") + or a.get("transition", "") + or ("硬切" if i < len(clips) - 1 else "结束"), "reference_image_index": ref_idx, } - voice = (xp.text_of(body, "voice_text") if body else "") + voice = xp.text_of(body, "voice_text") if body else "" if voice: voice_parts.append(voice) if not shot["scene_and_dialogue"]: @@ -883,30 +945,25 @@ def _script_from_xml(raw: str, job: ViralVideoJob) -> dict | None: base["title"] = base["overview"]["theme"] return base + def _step_script_generation(job: ViralVideoJob, intent: dict, image_analysis: dict) -> dict: """步骤 3: 编导分镜脚本生成(#2040:模板 + XML 解析;输出 copy_result 结构)。""" try: - from packages.shared.ai_service import call_llm from packages.application.viral_video.prompt_loader import ( get_template, render_user_prompt, ) from packages.application.viral_video.prompts import ( - FUSION_INSTRUCTIONS, GLOBAL_CONSTRAINTS, NEGATIVE_RULES, + FUSION_INSTRUCTIONS, + GLOBAL_CONSTRAINTS, + NEGATIVE_RULES, ) + from packages.shared.ai_service import call_llm except ImportError: return _fallback_script(job) products_summary = _build_products_summary(image_analysis) - style_hint = "无" - if isinstance(job.style_guide, dict): - style_hint = ( - f"节奏{job.style_guide.get('cut_speed', '')}、转场{job.style_guide.get('transition', '')}、" - f"色调{job.style_guide.get('color_grade', '')}、能量{job.style_guide.get('energy', '')}" - ) - dur = max(5, min(30, int(getattr(job, "duration", 15) or 15))) - ratio = getattr(job, "video_ratio", None) or "9:16" intent_str = "推广产品" key_msgs = "产品亮点" @@ -916,13 +973,6 @@ def _step_script_generation(job: ViralVideoJob, intent: dict, image_analysis: di key_msgs = "、".join(intent.get("key_messages") or []) or key_msgs tone = intent.get("tone") or tone - _vs = (job.viral_structure or "").strip() - if _vs: - viral_structure_block = f"【{_vs}】—— 请严格按照这个爆款结构的节奏/段落顺序编排镜头、台词和情绪节点(开场钩子、痛点、反转、案例、行动号召等按结构走),不要打乱顺序" - else: - viral_structure_block = "未指定(自由编排,但仍需有钩子开头+产品展示+行动号召的基本节奏)" - - persona_hint = _persona_style_hint(getattr(job, "persona_id", "")) fusion_level = getattr(job, "fusion_level", "ai_polish") or "ai_polish" fusion_instruction = FUSION_INSTRUCTIONS.get(fusion_level, FUSION_INSTRUCTIONS["ai_polish"]) @@ -939,15 +989,21 @@ def _step_script_generation(job: ViralVideoJob, intent: dict, image_analysis: di ) user = render_user_prompt( template, - duration=dur, image_count=len(job.images or []), + duration=dur, + image_count=len(job.images or []), fusion_result=fusion_brief, image_analysis=products_summary, ) def _try_gen(model: str, temp: float, max_tok: int, label: str, tmo: int = 25): logger.info("[爆款视频] 编导脚本生成 model=%s label=%s timeout=%d", model, label, tmo) - raw = call_llm([{"role": "system", "content": system_tpl}, {"role": "user", "content": user}], - temperature=temp, max_tokens=max_tok, model=model, timeout=tmo) + raw = call_llm( + [{"role": "system", "content": system_tpl}, {"role": "user", "content": user}], + temperature=temp, + max_tokens=max_tok, + model=model, + timeout=tmo, + ) if not raw: return None normalized = _script_from_xml(raw, job) @@ -967,7 +1023,13 @@ def _step_script_generation(job: ViralVideoJob, intent: dict, image_analysis: di fallback_marker = "我最近在用的好物" in voiceover has_typo_henjin = "很近" in json.dumps(normalized, ensure_ascii=False) is_fallback = fallback_marker or shots_cnt < 1 or len(voiceover) < 20 or has_typo_henjin - logger.info("[爆款视频] 编导脚本结果 label=%s voiceover_len=%d shots=%d fallback=%s", label, len(voiceover), shots_cnt, is_fallback) + logger.info( + "[爆款视频] 编导脚本结果 label=%s voiceover_len=%d shots=%d fallback=%s", + label, + len(voiceover), + shots_cnt, + is_fallback, + ) return None if is_fallback else normalized _s = get_shared_settings() @@ -992,6 +1054,7 @@ def _step_script_generation(job: ViralVideoJob, intent: dict, image_analysis: di logger.warning("[爆款视频] 编导脚本生成异常: %s,使用兜底脚本", e, exc_info=True) return _fallback_script(job) + def _step_review(job: ViralVideoJob, copy_result: dict) -> dict: """步骤 4: 合规审核(#2040:使用 Reviewer + review 模板,6 维度 + 自动重写 1 次)。 @@ -1000,7 +1063,11 @@ def _step_review(job: ViralVideoJob, copy_result: dict) -> dict: try: from packages.application.viral_video.reviewer import Reviewer from packages.application.viral_video.schemas import ( - FusionResult, IntentResult, CoreMessage, PersonalBrand, ScriptSegment, + CoreMessage, + FusionResult, + IntentResult, + PersonalBrand, + ScriptSegment, ) except ImportError as e: logger.warning("[爆款视频] reviewer 模块不可用,跳过审核: %s", e) @@ -1010,14 +1077,17 @@ def _step_review(job: ViralVideoJob, copy_result: dict) -> dict: title = (copy_result or {}).get("title") or (copy_result or {}).get("overview", {}).get("theme", "") intent_data = job.intent_result or {} - core_msgs = [CoreMessage(text=str(m), must_keep=True, confidence=0.9) for m in (intent_data.get("key_messages") or [])] + core_msgs = [ + CoreMessage(text=str(m), must_keep=True, confidence=0.9) for m in (intent_data.get("key_messages") or []) + ] brands: list[PersonalBrand] = [] brand_text = intent_data.get("brand_text") or intent_data.get("suggested_title") or "" if brand_text: brands.append(PersonalBrand(text=str(brand_text), category="brand")) intent_obj = IntentResult( intent_summary=intent_data.get("intent", "") or "推广产品", - core_messages=core_msgs, personal_brands=brands, + core_messages=core_msgs, + personal_brands=brands, ) fusion_obj = FusionResult( title=title or "", @@ -1041,7 +1111,11 @@ def _step_review(job: ViralVideoJob, copy_result: dict) -> dict: try: rewritten = reviewer.rewrite(fusion_obj, review_res, intent_obj, fusion_level) if rewritten and (rewritten.title or rewritten.script_segments): - new_voice = rewritten.script_segments[0].text if rewritten.script_segments else (rewritten.hook or voiceover) + new_voice = ( + rewritten.script_segments[0].text + if rewritten.script_segments + else (rewritten.hook or voiceover) + ) new_copy = dict(copy_result) new_copy["voiceover_script"] = new_voice new_copy["final_copy"] = new_voice @@ -1057,7 +1131,10 @@ def _step_review(job: ViralVideoJob, copy_result: dict) -> dict: "passed": review_res.passed, "score": 90 if review_res.passed else 60, "details": {i.dimension: i.text for i in review_res.issues}, - "issues": [{"dimension": i.dimension, "severity": i.severity, "location": i.location, "text": i.text} for i in review_res.issues], + "issues": [ + {"dimension": i.dimension, "severity": i.severity, "location": i.location, "text": i.text} + for i in review_res.issues + ], } if rewritten_voice is not None: result["rewritten_copy"] = new_copy @@ -1068,6 +1145,7 @@ def _step_review(job: ViralVideoJob, copy_result: dict) -> dict: logger.warning("[爆款视频] 审核异常,跳过: %s", e, exc_info=True) return {"passed": True, "score": 75, "details": {}, "issues": []} + def _step_tts(job: ViralVideoJob, voiceover_script: str): """步骤 5: CosyVoice 整段配音 → 返回本地 MP3 Path;失败返回 None。""" try: @@ -1106,6 +1184,7 @@ def _step_tts(job: ViralVideoJob, voiceover_script: str): logger.warning("[爆款视频] TTS 配音失败: %s", e, exc_info=True) return None + def _upload_tts_to_oss(job: ViralVideoJob, tts_path) -> str | None: """把 TTS 本地 mp3 上传到 OSS,返回公网 URL(供 Seedance 做 reference_audios 口型驱动用)。""" if tts_path is None: @@ -1125,6 +1204,7 @@ def _upload_tts_to_oss(job: ViralVideoJob, tts_path) -> str | None: logger.warning("[爆款视频] TTS 上传 OSS 失败: %s", e, exc_info=True) return None + def _assemble_seedance_prompt(copy_result: dict, job: ViralVideoJob) -> str: """把编导脚本拼成 Seedance 长 prompt。""" if not isinstance(copy_result, dict) or not copy_result: @@ -1175,6 +1255,7 @@ def _assemble_seedance_prompt(copy_result: dict, job: ViralVideoJob) -> str: lines.append(",".join([str(x) for x in np if x])) return "\n".join(lines) + def _step_render(job: ViralVideoJob, copy_result: dict, tts_audio_url: str | None) -> tuple[str, dict | None]: """步骤 6: v1.6 单次 Seedance 生成(不再分段/拼接)。 @@ -1288,6 +1369,7 @@ def _step_render(job: ViralVideoJob, copy_result: dict, tts_audio_url: str | Non logger.info("[爆款视频] 单次生成完成: path=%s size=%d usage=%s", video_path, size, usage) return str(video_path), (usage if isinstance(usage, dict) else None) + def _step_upload(job: ViralVideoJob, video_path: str) -> str: """步骤 7: OSS 上传。""" from video_processing.oss_helpers import upload_to_oss @@ -1300,6 +1382,7 @@ def _step_upload(job: ViralVideoJob, video_path: str) -> str: raise RuntimeError(f"OSS 上传失败: storage_key={storage_key}") return video_url + def _wait_oss_ready(url: str, timeout_sec: int = 10) -> bool: """轮询 OSS 公网 URL,直到 HEAD 返回 200 或超时。 用于缓解 OSS 上传后 1-5s 公网 eventual consistency 导致的 NoSuchKey。 @@ -1320,8 +1403,10 @@ def _wait_oss_ready(url: str, timeout_sec: int = 10) -> bool: logger.warning("[爆款视频] OSS 成片在 %ds 内未就绪 last_status=%s url=%s", timeout_sec, last_status, url[:120]) return False + # ── 主编排器 ──────────────────────────────────────────────────────────── + @shared_task( bind=True, max_retries=2, @@ -1406,6 +1491,7 @@ def run_viral_video_pipeline(self: Task, job_id: str) -> dict: 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: """旧 confirm-intent 路径兼容:从 WAIT_USER_CONFIRM 跑完整个渲染。""" @@ -1427,6 +1513,7 @@ def resume_viral_video_pipeline(self: Task, job_id: str) -> dict: 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: """独立的视频风格分析任务。""" @@ -1450,8 +1537,10 @@ def run_video_style_analysis(self: Task, job_id: str) -> dict: if session: session.close() + # ── 失败处理 ──────────────────────────────────────────────────────────── + def _mark_failed_and_notify(job_id: str, session, repo, job, err_msg: str, stage: str = "") -> None: """标记任务失败并通知。若传入的 session 已失效(因前面异常导致 rollback 状态), 会自动 fallback 到新建 SessionLocal 重新标记,确保状态一定落库。""" @@ -1492,8 +1581,10 @@ def _mark_failed_and_notify(job_id: str, session, repo, job, err_msg: str, stage event_type="viral_video:failed", ) + # ── v1.5/v1.6 三步分步流水线 Celery 任务 ───────────────────────────────── + @shared_task( bind=True, max_retries=1, @@ -1561,6 +1652,7 @@ def run_viral_video_analyze(self: Task, job_id: str) -> dict: if session: session.close() + @shared_task( bind=True, max_retries=1, @@ -1663,6 +1755,7 @@ def run_viral_video_generate_copy(self: Task, job_id: str) -> dict: if session: session.close() + def _quick_compliance_blacklist_check(copy_result: dict) -> None: """阶段2快速黑名单检查:不调用 LLM,只扫描高风险关键词;命中则在 voiceover 中就地替换。 @@ -1704,6 +1797,7 @@ def _quick_compliance_blacklist_check(copy_result: dict) -> None: for bk, bv in BLACKLIST.items(): copy_result[k] = copy_result[k].replace(bk, bv) + def _try_refund_viral_video(job: ViralVideoJob) -> None: """爆款视频生成失败:若已预扣积分则全额退款。""" try: @@ -1733,6 +1827,7 @@ def _try_refund_viral_video(job: ViralVideoJob) -> None: except Exception: logger.exception("[爆款视频] 失败退款异常 job_id=%s", job.id) + def _settle_viral_video(job: ViralVideoJob, usage: dict | None) -> None: """爆款视频生成成功:按实际 usage 结算,多退少补,写 credits_cost。""" try: @@ -1803,6 +1898,7 @@ def _settle_viral_video(job: ViralVideoJob, usage: dict | None) -> None: job.credits_cost = float(getattr(job, "credits_prepaid", 0) or 0) job.credits_prepaid = 0.0 + def _run_render_pipeline(job_id: str, session, repo, job) -> dict: """v1.6.1 阶段3:出片前合规审核(LLM 深度)→ TTS → Seedance → Upload → Completed。 @@ -1889,6 +1985,7 @@ def _run_render_pipeline(job_id: str, session, repo, job) -> dict: logger.info("[爆款视频] 任务完成: job_id=%s video_url=%s", job_id, video_url) return {"ok": True, "job_id": job_id, "video_url": video_url} + @shared_task(bind=True, max_retries=2, name="worker.run_viral_video_render") def run_viral_video_render(self: Task, job_id: str) -> dict: """v1.6 阶段3:TTS + 单次 Seedance 生成 + 上传。"""