Files
yingqing/core/backend/apps/ai/video_digest.py
T

573 lines
22 KiB
Python

"""上传视频提炼 —— 参考视频 → 可人工逐镜编辑的中文分镜稿。
优先把**完整视频(含音轨)**内联给 Gemini 3.1 Pro,跟官网「直接上传视频」同一条能力:
模型能看全片、听口播,不会只拿到稀疏静帧。文件太大塞不进请求时,才退回 ffmpeg 抽帧。
拆视频必须会看视频/图。默认语言模型现在可能是纯文本豆包,不能用 get_default_model(TEXT)。
固定钉 Gemini 3.1 Pro 官转(``gemini-3.1-pro-preview``,展示名带「官转」优先)。
"""
from __future__ import annotations
import base64
import json
import logging
import math
import re
import shutil
import subprocess
import tempfile
import uuid
from dataclasses import dataclass
from functools import lru_cache
from pathlib import Path
logger = logging.getLogger(__name__)
from django.conf import settings
# 上传限制:超了直接 400,不进 ffmpeg,也不花模型钱
ALLOWED_SUFFIXES = (".mp4", ".mov", ".m4v", ".webm")
MAX_UPLOAD_BYTES = 200 * 1024 * 1024 # 200 MB
MAX_DURATION_SECONDS = 180 # 3 分钟。带货参考片远短于此;更长的帧采样密度不够,拆出来也是错的
# 官网直接传视频的上限大约是请求 20MB;base64 会胀到 4/3,所以原文件卡在 15MB。
INLINE_VIDEO_MAX_BYTES = 15 * 1024 * 1024
_SUFFIX_MIME = {
".mp4": "video/mp4",
".m4v": "video/mp4",
".mov": "video/quicktime",
".webm": "video/webm",
}
# 抽帧只在整段视频塞不进请求时启用。
SECONDS_PER_FRAME = 2
MIN_FRAMES = 8
MAX_FRAMES = 36
FRAME_WIDTH = 768
FRAME_QUALITY = 3
DIGEST_MAX_TOKENS = 8192
_SHOT_MARK = re.compile(r"【第\s*\d+\s*镜】")
_FFMPEG_TIMEOUT = 60
class VideoDigestError(ValueError):
"""用户可见的失败(文件不合格 / ffmpeg 读不动),一律 400。"""
@dataclass(frozen=True)
class VideoFrame:
at_seconds: int
jpeg: bytes
def as_data_url(self) -> str:
return "data:image/jpeg;base64," + base64.b64encode(self.jpeg).decode("ascii")
@dataclass(frozen=True)
class DigestVideo:
mime: str
data: bytes
def as_data_url(self) -> str:
return f"data:{self.mime};base64," + base64.b64encode(self.data).decode("ascii")
# --------------------------------------------------------------------------- #
# skill 加载
# --------------------------------------------------------------------------- #
def _skill_dir() -> Path:
"""与 script_agent._skill_dir 同源:优先 BASE_DIR/skills(镜像内),回落仓库根(本地旧布局)。"""
base = Path(settings.BASE_DIR)
for cand in (base / "skills", base.parent.parent / "skills"):
if (cand / "video-shot-digest").is_dir():
return cand / "video-shot-digest"
return base / "skills" / "video-shot-digest"
@lru_cache(maxsize=1)
def load_digest_skill() -> str:
main = _skill_dir() / "SKILL.md"
if main.exists():
return main.read_text(encoding="utf-8")
# 兜底:skill 丢了也别整条链路挂掉,退化成一句话提示词(产出会明显变差,交接文档已注明须带 skills 目录)
return (
"你是分镜拆解 agent。输入是一条电商短视频的完整文件,含画面和口播。"
"必须覆盖全片,从 0 秒写到片尾,有几镜写几镜,不要概括成几大段。"
"每镜写画面和解说词,解说词按听到的口播逐字写。输出中文纯文本。"
)
# --------------------------------------------------------------------------- #
# ffmpeg:探时长 + 抽帧
# --------------------------------------------------------------------------- #
def _binary(name: str) -> str:
found = shutil.which(name)
if not found:
raise VideoDigestError("服务器暂时无法解析视频,请稍后再试")
return found
def probe_duration(path: str | Path) -> float:
"""ffprobe 读时长(秒)。读不到 = 不是能解的视频。"""
try:
out = subprocess.run(
[
_binary("ffprobe"), "-v", "error",
"-print_format", "json", "-show_format",
str(path),
],
capture_output=True, timeout=_FFMPEG_TIMEOUT, check=True,
).stdout
duration = float(json.loads(out)["format"]["duration"])
except VideoDigestError:
raise
except Exception as exc: # noqa: BLE001 — ffprobe 各种失败对用户是同一件事
raise VideoDigestError("这个视频读不出来,请换一个 mp4 / mov 文件") from exc
if duration <= 0:
raise VideoDigestError("这个视频读不出来,请换一个 mp4 / mov 文件")
return duration
def plan_frame_times(duration: float) -> list[int]:
"""均匀采样时间点。取每段的**中点**,避开首尾黑场与片尾卡片。"""
count = max(MIN_FRAMES, min(MAX_FRAMES, math.ceil(duration / SECONDS_PER_FRAME)))
step = duration / count
return [int(step * (i + 0.5)) for i in range(count)]
def extract_frames(path: str | Path, times: list[int]) -> list[VideoFrame]:
"""逐时间点抽一帧。``-ss`` 放在 ``-i`` 前走关键帧快速定位,每帧约几十毫秒。"""
ffmpeg = _binary("ffmpeg")
frames: list[VideoFrame] = []
for at in times:
try:
done = subprocess.run(
[
ffmpeg, "-v", "error", "-ss", str(at), "-i", str(path),
"-frames:v", "1", "-vf", f"scale={FRAME_WIDTH}:-2",
"-q:v", str(FRAME_QUALITY), "-f", "image2", "-",
],
capture_output=True, timeout=_FFMPEG_TIMEOUT, check=True,
)
except Exception: # noqa: BLE001 — 单帧抽失败(定位越界等)跳过,别拖垮整次提炼
continue
if done.stdout:
frames.append(VideoFrame(at_seconds=at, jpeg=done.stdout))
if not frames:
raise VideoDigestError("没能从这个视频里取到画面,请换一个文件")
return frames
def _write_upload(upload) -> tuple[str, str, int]:
"""校验后缀和体积,把上传落到临时文件。返回 (path, suffix, size)。调用方负责删除。"""
name = (getattr(upload, "name", "") or "").lower()
if not name.endswith(ALLOWED_SUFFIXES):
raise VideoDigestError("只支持 mp4 / mov / m4v / webm 四种视频格式")
size = getattr(upload, "size", 0) or 0
if size > MAX_UPLOAD_BYTES:
raise VideoDigestError(f"视频不能超过 {MAX_UPLOAD_BYTES // 1024 // 1024} MB,请压缩后再传")
suffix = Path(name).suffix or ".mp4"
tmp = tempfile.NamedTemporaryFile(suffix=suffix, delete=False)
try:
for chunk in upload.chunks():
tmp.write(chunk)
tmp.flush()
finally:
tmp.close()
return tmp.name, suffix, Path(tmp.name).stat().st_size
def _materialize_upload(upload) -> tuple[str, str, int, float]:
"""校验 → 落盘 → 探时长。返回 (path, suffix, size, duration),调用方负责删文件。"""
path, suffix, size = _write_upload(upload)
try:
duration = probe_duration(path)
except Exception:
Path(path).unlink(missing_ok=True)
raise
if duration > MAX_DURATION_SECONDS:
Path(path).unlink(missing_ok=True)
raise VideoDigestError(
f"视频不能超过 {MAX_DURATION_SECONDS // 60} 分钟,请剪出要参考的那一段再传"
)
return path, suffix, size, duration
def _compress_video(path: str) -> bytes | None:
"""压到能内联的体积。失败返回 None,由调用方改抽帧。"""
ffmpeg = _binary("ffmpeg")
out = tempfile.NamedTemporaryFile(suffix=".mp4", delete=False)
out.close()
try:
done = subprocess.run(
[
ffmpeg, "-v", "error", "-y", "-i", path,
"-c:v", "libx264", "-preset", "veryfast", "-crf", "28",
"-vf", "scale='min(1280,iw)':-2",
"-c:a", "aac", "-b:a", "64k",
"-movflags", "+faststart", out.name,
],
capture_output=True, timeout=_FFMPEG_TIMEOUT * 3,
)
if done.returncode != 0:
return None
data = Path(out.name).read_bytes()
if not data or len(data) > INLINE_VIDEO_MAX_BYTES:
return None
return data
except Exception: # noqa: BLE001 — 压缩失败就抽帧,别挡住提炼
return None
finally:
Path(out.name).unlink(missing_ok=True)
def _native_video(path: str, size: int, suffix: str) -> DigestVideo | None:
mime = _SUFFIX_MIME.get(suffix.lower(), "video/mp4")
if size <= INLINE_VIDEO_MAX_BYTES:
return DigestVideo(mime=mime, data=Path(path).read_bytes())
compressed = _compress_video(path)
if compressed:
return DigestVideo(mime="video/mp4", data=compressed)
return None
def digest_input_from_upload(upload) -> tuple[DigestVideo | None, list[VideoFrame], float]:
"""优先整段视频(含音轨);塞不进请求才抽帧。"""
path = ""
try:
path, suffix, size, duration = _materialize_upload(upload)
video = _native_video(path, size, suffix)
if video is not None:
return video, [], duration
return None, extract_frames(path, plan_frame_times(duration)), duration
finally:
if path:
Path(path).unlink(missing_ok=True)
def frames_from_upload(upload) -> tuple[list[VideoFrame], float]:
"""校验上传文件 → 落临时盘 → 探时长 → 抽帧。单测与抽帧兜底用。"""
path = ""
try:
path, _suffix, _size, duration = _materialize_upload(upload)
return extract_frames(path, plan_frame_times(duration)), duration
finally:
if path:
Path(path).unlink(missing_ok=True)
# --------------------------------------------------------------------------- #
# 组装多模态消息
# --------------------------------------------------------------------------- #
def build_digest_messages(
frames: list[VideoFrame] | None = None,
duration: float = 0,
*,
product_hint: str = "",
video: DigestVideo | None = None,
) -> list[dict]:
"""system = 拆解 skill;user = 完整视频(优先)或抽帧。"""
frames = frames or []
if video is not None:
head = [
f"这是一条时长约 {round(duration)} 秒的电商带货短视频的**完整文件**(含画面和口播音轨)。",
"请按技能还原全片分镜:从 0 秒写到片尾,有几镜写几镜。",
"解说词按听到的口播逐字写;画面上的花字一并写进画面。",
]
else:
head = [
f"这是一条时长约 {round(duration)} 秒的电商带货短视频,",
f"按时间顺序均匀抽了 {len(frames)} 帧。每帧图前面标了它在原片中的时间点。",
"请按技能还原**全片**分镜:从 0 秒写到片尾,有几镜写几镜。",
]
if product_hint:
head.append(f"用户接下来想用这条片子的结构去拍自己的商品:{product_hint}。")
content: list[dict] = [{"type": "text", "text": "".join(head)}]
if video is not None:
data_url = video.as_data_url()
# 官转(New-API)把 OpenAI image_url 的 data URI 转成 Gemini inline_data,
# mime 从 data:video/mp4 头读取,模型按整段视频+音轨理解,等同官网直接上传。
content.append({"type": "image_url", "image_url": {"url": data_url}})
else:
for frame in frames:
content.append({"type": "text", "text": f"[第 {frame.at_seconds} 秒]"})
content.append({"type": "image_url", "image_url": {"url": frame.as_data_url()}})
return [
{"role": "system", "content": load_digest_skill()},
{"role": "user", "content": content},
]
def min_digest_shots(duration: float, frame_count: int) -> int:
"""短片至少 3 镜;20 秒以上约每 6 秒一镜,且不超过抽到的帧数。"""
if duration < 20 and frame_count < 8:
return 3
by_time = max(4, int(duration // 6))
if frame_count:
return min(frame_count, by_time)
return by_time
def validate_digest_text(text: str, *, duration: float = 0, frame_count: int = 0) -> str:
"""模型偶尔吐空、只写片头、或概括成几大段。废稿不塞给用户。"""
cleaned = (text or "").strip()
if len(cleaned) < 80 or "【" not in cleaned:
raise ValueError("视频拆解结果不完整")
shots = _SHOT_MARK.findall(cleaned)
if not shots:
raise ValueError("视频拆解结果不完整")
if duration >= 20 or frame_count >= 8:
needed = min_digest_shots(duration, frame_count)
if len(shots) < needed:
raise ValueError("视频拆解镜头过少,请重试")
return cleaned
# 拆视频必须会看图。火山豆包直连读不了这组帧图,会在 ARK 上挂满 120s。
# 只认中转站的 Gemini 3.1 Pro;展示名带「官转」的优先。
DIGEST_VISION_MODEL_NAME = "gemini-3.1-pro-preview"
def _is_digest_vision_model(model) -> bool:
from apps.ai.services import OFFICIAL_DIRECT_PROVIDERS
provider_name = getattr(getattr(model, "provider", None), "name", "") or ""
if provider_name in OFFICIAL_DIRECT_PROVIDERS:
return False
blob = f"{model.name} {model.display_name}".lower()
return (
model.name == DIGEST_VISION_MODEL_NAME
or "gemini-3.1" in blob
or "gemini 3.1" in blob
)
def resolve_digest_model_config(preferred_id=None):
"""视频提炼用的多模态文本模型:Gemini 3.1 Pro 官转。找不到不回落默认语言模型。
前端可传 model_config_id(跟生成脚本同一套下拉)。只有会看图的 Gemini 3.1 才认,
选了豆包等纯文本模型仍钉回官转,避免拆帧直接失败。
"""
from apps.ai.models import ModelConfig
qs = (
ModelConfig.objects.select_related("provider")
.filter(
capability=ModelConfig.Capability.TEXT,
status=ModelConfig.Status.ACTIVE,
provider__status="active",
)
)
if preferred_id:
chosen = qs.filter(pk=preferred_id).first()
if chosen is not None and _is_digest_vision_model(chosen):
return chosen
def _blob(model) -> str:
return " ".join(
filter(
None,
[
model.name,
model.display_name,
getattr(model.provider, "name", ""),
getattr(model.provider, "display_name", ""),
],
)
)
ranked = []
for model in qs:
if not _is_digest_vision_model(model):
continue
blob = _blob(model)
# 官转 > 精确模型名 > 其它 Gemini 3.1
score = 0
if "官转" in blob:
score += 100
if model.name == DIGEST_VISION_MODEL_NAME:
score += 20
if getattr(model.provider, "name", "") == "yunqi_gemini":
score += 5
ranked.append((score, model.created_at, model))
if not ranked:
return None
ranked.sort(key=lambda item: (-item[0], item[1]))
return ranked[0][2]
# --------------------------------------------------------------------------- #
# 入口:一次真实的计费调用
# --------------------------------------------------------------------------- #
def digest_project_video(*, project, user, upload, model_config_id=None) -> dict:
"""上传视频 → 分镜稿。抽帧在建任务之前做,文件不合格不占积分。"""
product = getattr(project, "product", None)
return _digest_video(
team=project.team,
user=user,
upload=upload,
project=project,
product_hint=" · ".join(
filter(None, [getattr(product, "title", ""), getattr(product, "category", "")])
),
model_config_id=model_config_id,
)
def digest_team_video(*, team, user, upload, model_config_id=None) -> dict:
"""视频提炼页:不绑项目,计费挂当前团队。"""
return _digest_video(
team=team,
user=user,
upload=upload,
project=None,
product_hint="",
model_config_id=model_config_id,
)
def _digest_video(*, team, user, upload, project=None, product_hint="", model_config_id=None) -> dict:
"""上传视频 → 分镜稿。抽帧在建任务之前做,文件不合格不占积分。"""
from django.db import transaction
from django.utils import timezone
from apps.ai.models import AITask
from apps.ai.services import create_ai_task, execute_routed_text_request
from apps.billing.pricing import quote_video_digest
from apps.billing.services.ledger import charge_reserved_credit, reserve_credit
video, frames, duration = digest_input_from_upload(upload)
model_config = resolve_digest_model_config(preferred_id=model_config_id)
if model_config is None:
raise VideoDigestError("视频提炼需要 Gemini 3.1 Pro(会看图),当前没有启用,请联系管理员")
logger.info(
"video digest using %s:%s (%s) input=%s duration=%.1fs bytes=%s frames=%s",
model_config.provider.name,
model_config.name,
model_config.display_name,
"native_video" if video is not None else "frames",
duration,
len(video.data) if video is not None else 0,
len(frames),
)
messages = build_digest_messages(
frames, duration, product_hint=product_hint, video=video
)
request_payload = {
"model": model_config.name,
"endpoint": model_config.endpoint,
"feature": "video_remix" if project is None else "video_digest",
"duration_seconds": round(duration, 2),
"input": "native_video" if video is not None else "frames",
"frame_count": len(frames),
"video_bytes": len(video.data) if video is not None else 0,
"frame_times": [f.at_seconds for f in frames],
}
quote = quote_video_digest(team=team, model_config=model_config)
if quote.meta.get("rate"):
request_payload = {**request_payload, "points_per_yuan_snapshot": quote.meta["rate"]}
if project is not None:
task = create_ai_task(
project=project,
user=user,
task_type=AITask.Type.VIDEO_DIGEST,
model_config=model_config,
# 帧是几百 KB base64,绝不进 request_payload(会把 AITask 表撑爆),只记形状
request_payload=request_payload,
quote=quote,
)
else:
try:
with transaction.atomic():
task = AITask.objects.create(
team=team,
created_by=user,
project=None,
task_type=AITask.Type.VIDEO_DIGEST,
status=AITask.Status.CREATED,
model_config=model_config,
idempotency_key=f"video_digest:{team.id}:{uuid.uuid4()}",
request_payload=request_payload,
estimated_cost=quote.points,
base_cost=quote.base_cost_yuan,
)
reserve_credit(team=team, user=user, task=task, amount=quote.points)
task.status = AITask.Status.RESERVED
task.save(update_fields=["status", "updated_at"])
except ValueError as exc:
if "insufficient credit" in str(exc).lower():
raise VideoDigestError("团队余额不足,请充值后重试") from exc
raise VideoDigestError(str(exc)) from exc
reservation = task.credit_reservation
try:
task.status = AITask.Status.SUBMITTED
task.submitted_at = timezone.now()
task.save(update_fields=["status", "submitted_at", "updated_at"])
routed = execute_routed_text_request(
task=task,
primary_model=model_config,
messages=messages,
streaming=True,
structured_output=False,
business_operation="video_digest",
temperature=0.4,
validate_text=lambda text: validate_digest_text(
text, duration=duration, frame_count=len(frames) or (24 if video else 0)
),
extra_body={"max_tokens": DIGEST_MAX_TOKENS},
request_summary={
"duration_seconds": round(duration, 2),
"frame_count": len(frames),
"input": "native_video" if video is not None else "frames",
},
allow_retry=False,
allow_fallback=False,
)
_text, _response, digest = routed.value
except Exception as exc: # noqa: BLE001
_fail_digest_task(task, reservation, str(exc))
raise
with transaction.atomic():
task.status = AITask.Status.SUCCEEDED
task.response_payload = {"digest": digest[:32000]}
task.actual_cost = task.estimated_cost
task.completed_at = timezone.now()
task.save(update_fields=["status", "response_payload", "actual_cost", "completed_at", "updated_at"])
charge_reserved_credit(reservation=reservation, actual_amount=task.actual_cost)
return {
"text": digest,
"chars": len(digest),
"frames": len(frames),
"input": "native_video" if video is not None else "frames",
"duration": round(duration, 1),
"task_id": str(task.id),
"estimated_cost": str(task.estimated_cost),
}
def _fail_digest_task(task, reservation, message: str) -> None:
from django.utils import timezone
from apps.ai.models import AITask
from apps.billing.services.ledger import release_credit
try:
task.status = AITask.Status.FAILED
task.error_message = message[:2000]
task.completed_at = timezone.now()
task.save(update_fields=["status", "error_message", "completed_at", "updated_at"])
finally:
try:
release_credit(reservation=reservation, reason=message[:200])
except Exception: # noqa: BLE001
pass