Files
xiaoxia-saas/packages/shared/ai_client.py
T
xiaoxia 2a739dee17
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 1s
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (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 / PR Build Web 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
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 51s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 1m7s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 1m27s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m0s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 4m9s
CI/CD Pipeline / Validate - Style (pull_request) Failing after 4m41s
AI Code Review / AI Code Review (pull_request) Has been cancelled
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been cancelled
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 5m0s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 10m11s
CI/CD Pipeline / Unit Tests (pull_request) Successful in 18m38s
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / CI Gate (pull_request) Failing after 2s
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
fix(viral-video): #2134 generate-copy 提速 + 细粒度 phase/phase_message
问题7(generate-copy 提速,目标 30-40s):
- 新增 doubao_fast_model 配置(默认 doubao-1-5-pro-32k-250115),结构化输出任务(意图解析/编导脚本/合规审核)改用快模型,不再使用慢推理模型 doubao-seed-1-6
- call_llm 扩展支持 model/max_tokens/system_prompt 参数;chat_completion 同步支持 model 覆盖
- 编导脚本 temperature 0.8 + max_tokens 2500(从 4096 收紧);意图解析 max_tokens 800;审核 max_tokens 500
- _SCRIPT_GENERATION_PROMPT 精简冗余描述(前导说明和关键要求章节从 ~70 行压到 ~30 行),减少输入/输出 token
- 合规审核异步后置:阶段2 generate-copy 只做关键字黑名单快速检查(不调用 LLM),LLM 深度审核移到阶段3 confirm-copy TTS 之前执行,不再阻塞前端展示脚本
- 新增 _quick_compliance_blacklist_check 处理常见广告法绝对化用语

问题8(细粒度 phase + phase_message):
- ViralVideoJob 新增 phase_message 字段(中文提示文案,前端轮询直接展示)
- SQLAlchemy ViralVideoJobModel 新增 current_stage/phase_message 列(current_stage 原已有但未持久化更新)
- repo 层 _to_domain/save/update 同步处理新字段
- alembic 090 迁移:幂等 ADD COLUMN phase_message VARCHAR(500)
- 新增 _set_stage 辅助:统一设置 current_stage + phase_message + Redis 推送 + DB 持久化
- 所有 celery task(analyze/generate-copy/render/one-click/resume)在关键节点调用 _set_stage 持久化阶段信息
- 阶段文案:analyzing_images→正在分析商品特征 / parsing_intent→正在解析文案意图 / generating_script→正在编排分镜脚本 / reviewing→合规审核中 / tts→正在合成AI配音 / rendering→正在生成视频 / uploading→正在上传视频
- ViralVideoJobResponse schema + _to_response 增加 current_stage/phase_message,前端轮询 GET /{job_id} 直接拿到

配套:
- .env.example + render_env.sh SHARED_SECRETS 补 DOUBAO_FAST_MODEL
- tests _FakeSettings/_make_job 同步新增字段
- _run_render_pipeline 兜底补生成分支不再同步调用 _step_review(由出片前统一审核处理)
2026-10-02 11:01:01 +08:00

475 lines
19 KiB
Python
Executable File
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.
"""豆包大模型 API 客户端(共享层).
API 和 Worker 两边共用。基于火山引擎方舟平台的 OpenAI 兼容接口。
使用方式:
from packages.shared.ai_client import get_doubao_client
client = get_doubao_client()
if client.is_available:
result = client.chat_completion(messages=[...])
"""
from __future__ import annotations
import logging
import os
import time
import uuid
from typing import Any, Optional
import httpx
from packages.shared.config import get_shared_settings
logger = logging.getLogger(__name__)
class DoubaoClient:
"""豆包大模型 API 客户端.
封装 OpenAI 兼容的 Chat Completion 接口,支持自动重试。
未配置 API Key 时 is_available 为 False,调用方应降级处理。
"""
def __init__(self) -> None:
settings = get_shared_settings()
self.api_key: str = settings.doubao_api_key
self.model: str = settings.doubao_model
self.base_url: str = settings.doubao_base_url.rstrip("/")
self.timeout: int = settings.doubao_timeout
self.max_retries: int = settings.doubao_max_retries
self.vision_model: str = settings.doubao_vision_model
self.vision_lite_model: str = settings.doubao_vision_lite_model
self.fast_model: str = settings.doubao_fast_model
def embed_text(self, text: str, timeout: int | None = None) -> list[float] | None:
"""调用豆包文本 Embedding API,返回浮点向量;失败返回 None。"""
if not self.is_available or not text or not text.strip():
return None
url = f"{self.base_url}/embeddings"
headers = {
"Authorization": f"Bearer {self.api_key}",
"Content-Type": "application/json",
}
payload: dict[str, Any] = {
"model": getattr(self, "embedding_model", None) or "doubao-embedding-large-text-240915",
"input": text.strip(),
"encoding_format": "float",
}
req_timeout = timeout or self.timeout
last_error: Exception | None = None
for attempt in range(self.max_retries + 1):
try:
resp = httpx.post(url, headers=headers, json=payload, timeout=req_timeout)
resp.raise_for_status()
data = resp.json()
emb_list = data.get("data") or []
if emb_list and isinstance(emb_list, list):
vec = emb_list[0].get("embedding")
if isinstance(vec, list) and vec:
return [float(x) for x in vec]
logger.warning("embedding 返回结构异常: %s", str(data)[:200])
return None
except Exception as e:
last_error = e
if attempt < self.max_retries:
wait = 0.5 * (2**attempt)
logger.warning(
"豆包 Embedding 调用失败,%.1fs 后重试 (%d/%d): %s", wait, attempt + 1, self.max_retries + 1, e
)
time.sleep(wait)
logger.error("豆包 Embedding 调用最终失败: %s", last_error)
return None
@property
def is_available(self) -> bool:
"""是否可用(配置了 API Key)."""
return bool(self.api_key)
def chat_completion(
self,
messages: list[dict[str, str]],
temperature: float = 0.7,
max_tokens: int = 1024,
model: str | None = None,
) -> Optional[str]:
"""调用 Chat Completion 接口.
Args:
messages: 对话消息列表,[{"role": "user"/"system"/"assistant", "content": "..."}]
temperature: 采样温度,0-2,默认0.7
max_tokens: 最大生成token数,默认1024
Returns:
模型返回的文本内容,失败返回 None
"""
if not self.is_available:
return None
url = f"{self.base_url}/chat/completions"
headers = {
"Authorization": f"Bearer {self.api_key}",
"Content-Type": "application/json",
}
payload: dict[str, Any] = {
"model": model or self.model,
"messages": messages,
"temperature": temperature,
"max_tokens": max_tokens,
}
last_error: Optional[Exception] = None
for attempt in range(self.max_retries + 1):
try:
response = httpx.post(
url,
headers=headers,
json=payload,
timeout=self.timeout,
)
response.raise_for_status()
data = response.json()
content = data["choices"][0]["message"]["content"]
return content.strip()
except Exception as e:
last_error = e
if attempt < self.max_retries:
wait = 0.5 * (2**attempt)
logger.warning(
"豆包API调用失败,%.1fs后重试 (第%d/%d次): %s",
wait,
attempt + 1,
self.max_retries + 1,
e,
)
time.sleep(wait)
logger.error("豆包API调用最终失败: %s", last_error)
return None
def vision_completion(
self,
messages: list[dict],
images: list[str] | None = None,
max_tokens: int = 2048,
temperature: float = 0.3,
timeout: int | None = None,
model: str | None = None,
) -> Optional[str]:
"""调用豆包视觉理解 API(OpenAI 兼容多模态格式).
将 images 附加到最后一条 user message 的 content 中,
使用 vision_model(默认 doubao-1-5-vision-pro-250915)。
Args:
messages: 对话消息列表。最后一条 user message 会被注入图片内容。
images: 图片列表,支持 base64 data URI 或 HTTP(S) URL。
max_tokens: 最大生成 token 数,默认 2048。
temperature: 采样温度,默认 0.3(视觉任务偏低更稳定)。
timeout: 单次请求超时秒数,不传则使用默认 self.timeout。
Returns:
模型返回的文本内容,失败返回 None。
"""
if not self.is_available:
return None
# 构造多模态 content:先追加文本,再追加图片
vision_messages = []
for msg in messages:
vision_messages.append(dict(msg))
# 将图片注入最后一条 user message
if images and vision_messages:
# 找到最后一条 user message
for i in range(len(vision_messages) - 1, -1, -1):
if vision_messages[i].get("role") == "user":
text_content = vision_messages[i].get("content", "")
multi_content: list[dict[str, Any]] = []
if text_content:
multi_content.append({"type": "text", "text": text_content})
for img in images:
if img.startswith("data:") or img.startswith("http://") or img.startswith("https://"):
multi_content.append({"type": "image_url", "image_url": {"url": img}})
else:
# 当作 base64 编码
multi_content.append(
{"type": "image_url", "image_url": {"url": f"data:image/jpeg;base64,{img}"}}
)
vision_messages[i]["content"] = multi_content
break
url = f"{self.base_url}/chat/completions"
headers = {
"Authorization": f"Bearer {self.api_key}",
"Content-Type": "application/json",
}
payload: dict[str, Any] = {
"model": model or self.vision_model,
"messages": vision_messages,
"temperature": temperature,
"max_tokens": max_tokens,
}
req_timeout = timeout or self.timeout
last_error: Optional[Exception] = None
for attempt in range(self.max_retries + 1):
try:
response = httpx.post(
url,
headers=headers,
json=payload,
timeout=req_timeout,
)
response.raise_for_status()
data = response.json()
content = data["choices"][0]["message"]["content"]
return content.strip()
except Exception as e:
last_error = e
if attempt < self.max_retries:
wait = 0.5 * (2**attempt)
logger.warning(
"豆包视觉API调用失败,%.1fs后重试 (第%d/%d次): %s",
wait,
attempt + 1,
self.max_retries + 1,
e,
)
time.sleep(wait)
logger.error("豆包视觉API调用最终失败: %s", last_error)
return None
# ── 视频生成(Seedance 2.5,异步任务)────────────────────────────
def video_generation(
self,
prompt: str,
*,
image_url: str | None = None,
duration: int = 5,
ratio: str | None = "9:16",
resolution: str = "720p",
generate_audio: bool = True,
watermark: bool = False,
output_dir: str | None = None,
model: str | None = None,
reference_images: list[str] | None = None,
reference_audios: list[str] | None = None,
reference_videos: list[str] | None = None,
) -> str | None:
"""调用 Seedance 2.5 文生/图生视频(异步任务→轮询→下载),返回本地 MP4 路径;失败返回 None。
v1.6: 支持多参考图(产品素材)+ 参考音频(TTS口型驱动)+ 参考视频,单次生成最长 30 秒。
Args:
prompt: 文本提示词(含完整编导脚本:总览+场景光线+逐镜头时间轴+硬约束+负面词)
image_url: 首帧参考图 URL(可选,提供则走图生视频首帧模式,ratio 跟随首帧)
duration: 视频时长 4~30 秒
ratio: 宽高比 16:9/9:16/1:1/4:3/3:4/21:9/adaptive;image_url 存在时自动忽略
resolution: 480p/720p/1080p
generate_audio: 是否让模型原生合成音效/BGM(v1.6 默认 True,配合 reference_audios 做口型驱动)
watermark: 是否加水印
output_dir: 下载目录,默认 /tmp
model: 指定模型 ID;空则用 settings.doubao_video_model
reference_images: 多参考图 URL 列表(产品素材,最多30张;注意 image_url 为首帧单独传)
reference_audios: 参考音频 URL 列表(TTS口播,驱动口型,最多10段)
reference_videos: 参考视频 URL 列表(风格参考)
Returns:
本地 MP4 文件路径,失败返回 None。
"""
if not self.is_available:
return None
if not prompt or not prompt.strip():
return None
settings = get_shared_settings()
poll_interval = getattr(settings, "doubao_video_poll_interval", 10) or 10
total_timeout = getattr(settings, "doubao_video_timeout", 900) or 900
default_video_model = getattr(settings, "doubao_video_model", None) or "doubao-seedance-2-5-260628"
video_model = model or default_video_model
content: list[dict[str, Any]] = [{"type": "text", "text": prompt.strip()}]
if image_url:
content.append({"type": "image_url", "image_url": {"url": image_url}})
# v1.6: 多参考图(产品素材)
if reference_images:
for url in reference_images[:30]:
if url and isinstance(url, str):
content.append({"type": "image_url", "image_url": {"url": url}})
create_payload: dict[str, Any] = {
"model": video_model,
"content": content,
"generate_audio": bool(generate_audio),
"duration": int(duration),
"resolution": resolution,
"watermark": bool(watermark),
}
# v1.6: 参考音频(TTS 驱动口型)
if reference_audios:
create_payload["reference_audios"] = [
{"url": u, "role": "audio_url"} for u in reference_audios[:10] if u and isinstance(u, str)
]
# v1.6: 参考视频(风格参考)
if reference_videos:
create_payload["reference_videos"] = [{"url": u} for u in reference_videos[:5] if u and isinstance(u, str)]
# Bug #2110 / v1.6: ratio=None 时不传(首帧图生视频跟随原图比例)
if ratio and not image_url:
create_payload["ratio"] = ratio
headers = {
"Authorization": f"Bearer {self.api_key}",
"Content-Type": "application/json",
}
create_url = f"{self.base_url}/contents/generations/tasks"
logger.info(
"Seedance 创建任务请求: url=%s model=%s duration=%ds ratio=%s gen_audio=%s image_url=%s ref_imgs=%d ref_audios=%d ref_videos=%d",
create_url,
video_model,
duration,
ratio or "(follow-image)",
generate_audio,
bool(image_url),
len(reference_images or []),
len(reference_audios or []),
len(reference_videos or []),
)
# 1) 创建任务(带重试)
task_id: str | None = None
last_error: Exception | None = None
for attempt in range(self.max_retries + 1):
try:
resp = httpx.post(create_url, headers=headers, json=create_payload, timeout=self.timeout)
# 测试环境下 MagicMock().status_code 是 MagicMock,与 int 比较会抛 TypeError;
# 用显式 int() 转换+类型判断,避免误判。
try:
_status = int(resp.status_code)
except (TypeError, ValueError):
_status = 200
if _status >= 400:
# 把响应体完整打出来(通常含 error.code/message,能直接定位:模型未开通/Key 无权限/模型 ID 错误)
logger.error(
"Seedance 创建任务 HTTP %d: body=%s",
resp.status_code,
(resp.text or "")[:1000],
)
resp.raise_for_status()
data = resp.json()
task_id = data.get("id")
if task_id:
break
last_error = RuntimeError(f"create task returned no id: {str(data)[:200]}")
except Exception as e:
last_error = e
if attempt < self.max_retries:
wait = 0.5 * (2**attempt)
logger.warning(
"Seedance 创建任务失败,%.1fs 后重试 (%d/%d): %s", wait, attempt + 1, self.max_retries + 1, e
)
time.sleep(wait)
if not task_id:
logger.error(
"Seedance 创建任务最终失败: model=%s base_url=%s err=%s 【排查建议】"
"1) 确认方舟控制台已开通 Doubao-Seedance-2.5 模型;"
"2) DOUBAO_API_KEY 对应的账号有该模型调用权限;"
"3) DOUBAO_BASE_URL 必须为 https://ark.cn-beijing.volces.com/api/v3;"
"4) 若控制台用「推理接入点」(endpoint),请把 DOUBAO_VIDEO_MODEL 改为 ep-xxx 接入点 ID。",
video_model,
self.base_url,
last_error,
)
return None
logger.info(
"Seedance 任务已创建: task_id=%s model=%s duration=%ds gen_audio=%s",
task_id,
video_model,
duration,
generate_audio,
)
# 2) 轮询状态
poll_url = f"{create_url}/{task_id}"
deadline = time.time() + total_timeout
video_url: str | None = None
last_status: str = "queued"
while time.time() < deadline:
try:
resp = httpx.get(poll_url, headers=headers, timeout=self.timeout)
try:
if int(getattr(resp, "status_code", 200)) >= 400:
resp.raise_for_status()
except (TypeError, ValueError):
pass
data = resp.json()
status = data.get("status", "")
last_status = status
if status == "succeeded":
content_obj = data.get("content") or {}
video_url = content_obj.get("video_url")
if video_url:
break
last_error = RuntimeError(f"task succeeded but no video_url: {str(data)[:300]}")
break
if status == "failed":
err = data.get("error") or {}
last_error = RuntimeError(f"task failed: {err.get('code','')} {err.get('message','')}")
break
if status in ("expired", "cancelled"):
last_error = RuntimeError(f"task {status}")
break
# queued / running: 继续轮询
except httpx.HTTPStatusError as e:
last_error = e
logger.warning(
"Seedance 轮询 HTTP %d: body=%s",
e.response.status_code,
(e.response.text or "")[:500],
)
except Exception as e:
last_error = e
logger.debug("Seedance 轮询异常: %s", e)
time.sleep(poll_interval)
if not video_url:
logger.error("Seedance 任务未成功: task_id=%s status=%s err=%s", task_id, last_status, last_error)
return None
# 3) 下载到本地
try:
out_dir = output_dir or "/tmp"
os.makedirs(out_dir, exist_ok=True)
local_path = f"{out_dir}/seedance_{task_id}_{uuid.uuid4().hex[:8]}.mp4"
with httpx.stream("GET", video_url, timeout=300) as r:
r.raise_for_status()
with open(local_path, "wb") as f:
for chunk in r.iter_bytes(chunk_size=1024 * 256):
if chunk:
f.write(chunk)
logger.info("Seedance 视频下载完成: %s (%d bytes)", local_path, os.path.getsize(local_path))
return local_path
except Exception as e:
logger.error("Seedance 视频下载失败: %s", e)
return None
# ── 单例 ─────────────────────────────────────────────────────────────────────
_client: Optional[DoubaoClient] = None
def get_doubao_client() -> DoubaoClient:
"""获取豆包客户端单例."""
global _client
if _client is None:
_client = DoubaoClient()
return _client