917 lines
38 KiB
Python
917 lines
38 KiB
Python
"""视频复刻:参考视频 + 商品图/人物图 → Seedance 换商品或换角色。
|
||
|
||
不新建任务类型、不接检测/抠图。提交仍走 submit_free_video,只把
|
||
feature=video_replace 和 replace_mode 写进 payload,提示词由后端写死。
|
||
|
||
真人参考必须先送火山素材库审核,过审后用 asset:// 生成;审核失败不扣费。
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
import uuid
|
||
from decimal import Decimal
|
||
|
||
from django.conf import settings
|
||
from django.db import transaction
|
||
from django.db.models import Q
|
||
|
||
from apps.assets.models import Asset, Model
|
||
from apps.billing.models import CreditAccount
|
||
from apps.billing.pricing import quote_video_estimate, video_reserve_amount
|
||
from apps.products.models import Product
|
||
|
||
from .free_video import (
|
||
FREE_VIDEO_MODELS,
|
||
HIGH_RES_MODEL,
|
||
REPLACE_MODEL,
|
||
IN_FLIGHT_STATUSES,
|
||
RATIOS,
|
||
RESOLUTIONS,
|
||
_reap_stale_free_video_tasks,
|
||
model_duration_range,
|
||
serialize_free_video_task,
|
||
start_pending_free_video,
|
||
submit_free_video,
|
||
)
|
||
from .models import AITask, ModelConfig
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
FEATURE = "video_replace"
|
||
REPLACE_MODES = {"product", "character"}
|
||
MAX_IMAGES = 9
|
||
LEGACY_PROMPT_PREFIX = "[视频复刻]"
|
||
# 商品三视图在商品库是 standalone Asset(metadata.view=three_view + product_id),
|
||
# 和项目内的 BaseAssetGroup 是两套存储,这里直接取商品库那份。
|
||
TRIVIEW_LABEL = "目标商品三视图"
|
||
# 复刻统一走 Seedance 2.5,单次能出 30 秒,所以参考视频也按 30 秒收口:
|
||
# 商品模式的参考视频只喂提炼模型,角色模式直传火山,两边都不需要超过成片上限的素材。
|
||
# 半秒容差同 media_probe 的老规矩(30 秒成片常被探成 30.02)。
|
||
REPLACE_REF_DURATION_MAX = 30.5
|
||
REVIEW_UNAVAILABLE = "素材审核服务暂不可用,请稍后重试"
|
||
REVIEW_FAILED = "参考素材未通过真人合规审核,请更换视频或图片后重试"
|
||
REVIEW_SUBMIT_FAILED = "素材提交审核失败,请稍后重试"
|
||
|
||
|
||
class VideoReplaceInProgress(ValueError):
|
||
"""本团队已有复刻在跑。视图转 409,并把在跑的那条带回去让前端接上。"""
|
||
|
||
def __init__(self, task):
|
||
super().__init__("已有一个视频正在复刻中,完成或取消后才能再提交")
|
||
self.task = task
|
||
|
||
|
||
def get_inflight_video_replace(team):
|
||
"""本团队正在跑的复刻任务(含审核中 / 提炼中)。没有返回 None。"""
|
||
_reap_stale_free_video_tasks(team=team)
|
||
return (
|
||
AITask.objects.filter(
|
||
team=team,
|
||
task_type=AITask.Type.FREE_VIDEO,
|
||
status__in=IN_FLIGHT_STATUSES,
|
||
is_deleted=False,
|
||
purged_at__isnull=True,
|
||
)
|
||
.filter(video_replace_q())
|
||
.order_by("-created_at")
|
||
.first()
|
||
)
|
||
|
||
# 商品复刻不再把参考视频直接丢给火山(实测换不动商品,出片跑偏),改走两步:
|
||
# 先用「提炼提示词」那一整套(同一份 SKILL.md + 同一个 Gemini 3.1 Pro)把参考视频拆成
|
||
# 中文分镜稿,再把分镜稿当导演脚本、连同商品图一起交给 Seedance。
|
||
# 角色复刻仍是「参考视频 + 角色图」直传,已验证有效,不要顺手一起改。
|
||
PRODUCT_DIGEST_HEAD = (
|
||
"下面是参考视频的完整分镜稿,它逐镜记录了这条片子的时间、景别、机位、运镜、画面、"
|
||
"人物动作与表情、台词旁白和声音。请把它当作导演脚本,按镜头顺序还原成一条新视频。"
|
||
)
|
||
PRODUCT_DIGEST_TAIL = (
|
||
"【复刻要求】\n"
|
||
"1. 严格按上面分镜稿的镜头顺序、时间分配、景别、机位、运镜和剪辑节奏还原全片,不要自行增删镜头。\n"
|
||
"2. 分镜稿里出现的原商品,全部替换成 @目标商品;商品的外形、颜色、材质、包装和标识必须与参考图完全一致,"
|
||
"不要改造型、不要改配色、不要凭空补细节。\n"
|
||
"{triview}"
|
||
"3. 人物、场景、光线、色调、动作和表情按分镜稿保持不变,只换商品。\n"
|
||
" · 全片的人物必须自始至终是同一个人:五官、发型、体态和服装在所有镜头里保持一致,"
|
||
"不要中途换脸、换人或换装。\n"
|
||
" · 参考图里如果有人物(例如服饰类的模特图、真人手持商品的图),画面中的人物就以那张参考图为准,"
|
||
"五官、发型、体态和服装必须与之一致;参考图里只有商品、没有人物时,就按分镜稿描述的人物去演,"
|
||
"不要从商品图上凭空套一张脸。\n"
|
||
"4. 台词/旁白按分镜稿逐字念出;原文提到旧商品名称的地方,改说「{subject}」。\n"
|
||
"5. 分镜稿里的「字幕」栏只是对原片的记录,不要把这些文字画到画面上。"
|
||
"{condense}"
|
||
)
|
||
|
||
# 参考视频比成片长时才追加(火山单次出片上限 15 秒,60 秒参考稿必须压缩改编,
|
||
# 否则模型会照着分镜稿从头拍、拍到 15 秒戛然而止)。
|
||
PRODUCT_CONDENSE_NOTE = (
|
||
"\n6. 参考视频约 {source} 秒,而本次成片只有 {output} 秒:请按分镜稿的结构和节奏**压缩改编**,"
|
||
"保留主线、最关键的卖点镜头和结尾收束,可以合并或省略次要镜头;"
|
||
"不要只拍开头几镜就停,也不要把画面加速成快放。"
|
||
)
|
||
# 商品库里有三视图时才追加这一条(没有就不提,免得模型去找一张不存在的图)。
|
||
PRODUCT_TRIVIEW_NOTE = (
|
||
" · @目标商品三视图 是这件商品的白底多角度图,用它确认商品的立体结构、比例和各个面的细节;"
|
||
"镜头转到任何角度,商品都要和三视图对得上。\n"
|
||
)
|
||
|
||
# 拆解期间的占位提示词。真正的提示词在 worker 拆完后回写。
|
||
DIGEST_PENDING_PROMPT = "正在提炼参考视频的分镜稿…"
|
||
|
||
# 分镜稿是给人逐镜校对的导演稿,整篇塞给视频模型太长也太杂,按需瘦身:
|
||
# ·「字幕」栏是原片花字的记录,留着会诱导模型把字画到画面上,与全站「成片无字幕」冲突;
|
||
# ·【拆解存疑】是给用户校对用的,视频模型用不上。
|
||
DIGEST_DROP_FIELDS = ("字幕",)
|
||
# 瘦身完还超长时再砍这几栏(信息量最低,砍掉不影响镜头结构)。
|
||
DIGEST_TRIM_FIELDS = ("音效", "背景音乐", "备注")
|
||
DIGEST_TAIL_SECTION = "【拆解存疑】"
|
||
MAX_DIGEST_CHARS = 4000
|
||
CHARACTER_PROMPT = (
|
||
"使用@参考视频作为镜头、节奏与口播氛围基准,"
|
||
"将画面中需要替换的原人物完整替换为@目标角色。"
|
||
"保留参考视频的商品、场景、镜头运动、剪辑节奏与口播氛围,"
|
||
"角色五官、发型、体态必须与参考图一致,不要改变原片构图和商品展示。"
|
||
)
|
||
|
||
|
||
def _drop_digest_fields(lines: list[str], fields: tuple[str, ...]) -> list[str]:
|
||
prefixes = tuple(f"{name}:" for name in fields)
|
||
return [line for line in lines if not line.strip().startswith(prefixes)]
|
||
|
||
|
||
def compact_digest_for_video(digest_text: str) -> str:
|
||
"""分镜稿 → 交给视频模型的精简版。保留镜头结构,砍掉给人看的部分。"""
|
||
body = (digest_text or "").strip()
|
||
head, sep, _tail = body.partition(DIGEST_TAIL_SECTION)
|
||
if sep:
|
||
body = head.strip()
|
||
lines = _drop_digest_fields(body.split("\n"), DIGEST_DROP_FIELDS)
|
||
if sum(len(line) for line in lines) > MAX_DIGEST_CHARS:
|
||
lines = _drop_digest_fields(lines, DIGEST_TRIM_FIELDS)
|
||
return "\n".join(line for line in lines if line.strip()).strip()
|
||
|
||
|
||
def build_product_replace_prompt(
|
||
digest_text: str,
|
||
subject_name: str,
|
||
*,
|
||
has_triview: bool = False,
|
||
source_seconds: float = 0,
|
||
output_seconds: int = 0,
|
||
) -> str:
|
||
"""分镜稿 + 商品替换要求 → 交给 Seedance 的完整提示词。
|
||
|
||
@目标商品 必须与 references 里第一张商品图的 label 对上,否则 build_content_items
|
||
不会把它换成火山认的「图片N」指代。
|
||
"""
|
||
from .services import enforce_no_embedded_captions
|
||
|
||
body = compact_digest_for_video(digest_text)
|
||
if not body:
|
||
raise ValueError("参考视频拆解结果为空,请重试")
|
||
subject = (subject_name or "").strip() or "目标商品"
|
||
source = int(round(float(source_seconds or 0)))
|
||
output = int(output_seconds or 0)
|
||
# 差 3 秒以内不啰嗦,模型自己会收;差得多才明确要求压缩改编。
|
||
condense = (
|
||
PRODUCT_CONDENSE_NOTE.format(source=source, output=output)
|
||
if source and output and source - output >= 3
|
||
else ""
|
||
)
|
||
tail = PRODUCT_DIGEST_TAIL.format(
|
||
subject=subject,
|
||
triview=PRODUCT_TRIVIEW_NOTE if has_triview else "",
|
||
condense=condense,
|
||
)
|
||
return enforce_no_embedded_captions(f"{PRODUCT_DIGEST_HEAD}\n\n{body}\n\n{tail}")
|
||
|
||
|
||
def is_video_replace_task(task) -> bool:
|
||
payload = task.request_payload or {}
|
||
if payload.get("feature") == FEATURE:
|
||
return True
|
||
return str(payload.get("prompt") or "").startswith(LEGACY_PROMPT_PREFIX)
|
||
|
||
|
||
def video_replace_q() -> Q:
|
||
return Q(request_payload__feature=FEATURE) | Q(request_payload__prompt__startswith=LEGACY_PROMPT_PREFIX)
|
||
|
||
|
||
def serialize_video_replace_task(task, *, include_deleted_assets: bool = False) -> dict:
|
||
data = serialize_free_video_task(task, include_deleted_assets=include_deleted_assets)
|
||
payload = task.request_payload or {}
|
||
replace_mode = payload.get("replace_mode") or _legacy_replace_mode(payload.get("prompt") or "")
|
||
created = task.status == AITask.Status.CREATED
|
||
digesting = created and bool(payload.get("digest_pending"))
|
||
reviewing = created and bool(payload.get("review_pending")) and not digesting
|
||
data.update({
|
||
"feature": FEATURE,
|
||
"replace_mode": replace_mode,
|
||
"subject_name": payload.get("subject_name") or "",
|
||
"subject_source": payload.get("subject_source") or "",
|
||
"product_id": payload.get("product_id") or "",
|
||
"model_id": payload.get("model_id") or "",
|
||
"review_stage": "reviewing" if reviewing else "",
|
||
# 商品复刻专有:参考视频拆出来的分镜稿(给用户看,也便于排查出片跑偏)
|
||
"digest_stage": "digesting" if digesting else "",
|
||
"digest_text": str(payload.get("digest_text") or ""),
|
||
"digest_shots": payload.get("digest_shots") or 0,
|
||
"digest_source_name": payload.get("digest_source_name") or "",
|
||
"digest_source": payload.get("digest_source_ref") or None,
|
||
})
|
||
return data
|
||
|
||
|
||
def submit_video_replace(*, team, user, params: dict):
|
||
"""校验素材 → 送审 → 已过审则直接生成,否则建 CREATED 任务等绿盾。失败抛 ValueError。"""
|
||
replace_mode = str(params.get("replace_mode") or "").strip()
|
||
if replace_mode not in REPLACE_MODES:
|
||
raise ValueError("请选择替换商品或替换角色")
|
||
|
||
product_id = _optional_uuid(params.get("product_id"), "商品")
|
||
model_id = _optional_uuid(params.get("model_id"), "角色")
|
||
image_ids = _uuid_list(params.get("image_asset_ids"), "参考图")
|
||
has_product = product_id is not None
|
||
has_model = model_id is not None
|
||
has_temp = bool(image_ids)
|
||
|
||
if replace_mode == "product":
|
||
if has_model:
|
||
raise ValueError("商品复刻请选择商品,不要同时选择角色")
|
||
if has_product and has_temp:
|
||
raise ValueError("请从商品库选择,或临时上传商品图,不要混用")
|
||
if not has_product and not has_temp:
|
||
raise ValueError("请选择商品或上传商品参考图")
|
||
else:
|
||
if has_product:
|
||
raise ValueError("角色复刻请选择角色,不要同时选择商品")
|
||
if has_model and has_temp:
|
||
raise ValueError("请从人物库选择,或临时上传角色图,不要混用")
|
||
if not has_model and not has_temp:
|
||
raise ValueError("请选择角色或上传角色参考图")
|
||
|
||
# 单飞闸:一个团队同时只允许一条复刻在跑。刷新页面时前端拉状态有空档,
|
||
# 用户容易以为「没任务」再点一次,重复扣费。_reap_stale_free_video_tasks 会先回收
|
||
# 死任务(CREATED 超 16 分钟等),所以不会被永久锁住。
|
||
running = get_inflight_video_replace(team)
|
||
if running is not None:
|
||
raise VideoReplaceInProgress(running)
|
||
|
||
video = _team_asset(team, params.get("video_asset_id"), kind=Asset.Type.VIDEO, label="参考视频")
|
||
video_seconds = _asset_duration_seconds(video)
|
||
if video_seconds > REPLACE_REF_DURATION_MAX:
|
||
raise ValueError("参考视频不能超过 30 秒,请剪短后重试")
|
||
|
||
if has_product:
|
||
subject_name, image_refs, subject_source = _product_library_refs(team, product_id)
|
||
elif has_model:
|
||
subject_name, image_refs, subject_source = _character_library_refs(team, model_id)
|
||
else:
|
||
noun = "商品" if replace_mode == "product" else "角色"
|
||
subject_name, image_refs, subject_source = _temporary_image_refs(team, image_ids, noun=noun)
|
||
|
||
duration = _output_duration(params.get("duration"), video_seconds)
|
||
extra = {
|
||
"replace_mode": replace_mode,
|
||
"subject_name": subject_name,
|
||
"subject_source": subject_source,
|
||
"product_id": str(product_id) if product_id else "",
|
||
"model_id": str(model_id) if model_id else "",
|
||
}
|
||
base_params = {
|
||
"mode": "universal",
|
||
# 商品和角色都固定走 Seedance 2.5:只有它支持 30 秒单次出片,不接受前端指定别的档
|
||
"model": REPLACE_MODEL,
|
||
"aspect_ratio": str(params.get("aspect_ratio") or "9:16"),
|
||
"resolution": str(params.get("resolution") or "720p"),
|
||
"duration": duration,
|
||
"seed": params.get("seed", -1),
|
||
"generate_audio": True,
|
||
"feature": FEATURE,
|
||
}
|
||
|
||
if replace_mode == "product":
|
||
# 商品复刻:参考视频只喂给提炼模型,不进火山 references(所以参考视频本身不再送火山审核)。
|
||
# 两个「一定会失败」的前提在这里就查掉,别让用户白等一轮提炼才看到报错:
|
||
from .video_digest import resolve_digest_model_config
|
||
|
||
if resolve_digest_model_config() is None:
|
||
raise ValueError("视频提炼模型未配置,请联系管理员")
|
||
# 商品图仍要过审才能当生成参考。提前送审 → 审核和提炼并行跑,明确不过审的直接 400。
|
||
review_state = _ensure_replace_refs_reviewed(team, image_refs)
|
||
if review_state == "failed":
|
||
raise ValueError(REVIEW_FAILED)
|
||
image_refs = _refresh_replace_refs(team, image_refs)
|
||
# 先秒回一个 CREATED 任务,拆解这种慢活(Gemini 半分钟起)交给 worker。
|
||
extra.update({
|
||
"digest_source_asset_id": str(video.id),
|
||
"digest_source_name": video.name or "参考视频",
|
||
"digest_source_duration": round(video_seconds, 2),
|
||
# 参考视频虽然不进 references,历史卡的「原视频」和「重新生成」仍要拿到它,
|
||
# 这里存一份引用快照,免得序列化历史时每条再查一次库。
|
||
"digest_source_ref": _owned_ref(video, kind="video", role="reference_video", label="参考视频"),
|
||
})
|
||
return _create_reviewing_task(
|
||
team=team,
|
||
user=user,
|
||
params={
|
||
**base_params,
|
||
"prompt": DIGEST_PENDING_PROMPT,
|
||
"references": image_refs,
|
||
"extra_payload": extra,
|
||
},
|
||
digest_pending=True,
|
||
)
|
||
|
||
references = [
|
||
_owned_ref(video, kind="video", role="reference_video", label="参考视频"),
|
||
*image_refs,
|
||
]
|
||
review_state = _ensure_replace_refs_reviewed(team, references)
|
||
if review_state == "failed":
|
||
raise ValueError(REVIEW_FAILED)
|
||
references = _refresh_replace_refs(team, references)
|
||
extra["review_pending"] = review_state != "ready"
|
||
submit_params = {
|
||
**base_params,
|
||
"prompt": CHARACTER_PROMPT,
|
||
"references": references,
|
||
"extra_payload": extra,
|
||
}
|
||
if review_state == "ready":
|
||
_assert_replace_refs_ready(references)
|
||
return submit_free_video(team=team, user=user, params=submit_params)
|
||
return _create_reviewing_task(team=team, user=user, params=submit_params)
|
||
|
||
|
||
def advance_video_replace(task):
|
||
"""轮询审核中的复刻任务:失败则结束(不扣费),过审则预留积分并提交火山。"""
|
||
if not is_video_replace_task(task):
|
||
return task
|
||
if task.status != AITask.Status.CREATED:
|
||
from .free_video import finalize_free_video
|
||
|
||
return finalize_free_video(task=task)
|
||
|
||
payload = task.request_payload or {}
|
||
if payload.get("digest_pending"):
|
||
# worker 还在拆参考视频,提示词都没生成,别当成「审核中」反复送审。
|
||
# worker 挂了也不会永远卡着:CREATED 超 16 分钟由 _reap_stale_free_video_tasks 回收退费。
|
||
return task
|
||
references = list(payload.get("references") or [])
|
||
try:
|
||
state = _ensure_replace_refs_reviewed(task.team, references)
|
||
except ValueError as exc:
|
||
return _fail_reviewing_task(task, str(exc))
|
||
if state == "failed":
|
||
return _fail_reviewing_task(task, REVIEW_FAILED)
|
||
if state != "ready":
|
||
return task
|
||
|
||
refreshed = _refresh_replace_refs(task.team, references)
|
||
try:
|
||
_assert_replace_refs_ready(refreshed)
|
||
except ValueError as exc:
|
||
return _fail_reviewing_task(task, str(exc))
|
||
with transaction.atomic():
|
||
locked = AITask.objects.select_for_update().get(id=task.id)
|
||
if locked.status != AITask.Status.CREATED:
|
||
return locked
|
||
next_payload = dict(locked.request_payload or {})
|
||
next_payload["references"] = refreshed
|
||
next_payload["review_pending"] = False
|
||
locked.request_payload = next_payload
|
||
locked.save(update_fields=["request_payload", "updated_at"])
|
||
task = locked
|
||
return start_pending_free_video(task)
|
||
|
||
|
||
def run_replace_digest(task) -> None:
|
||
"""Worker:商品复刻第一道工序——参考视频 → 分镜稿 → 完整提示词 → 走审核提交火山。
|
||
|
||
这一步失败 = 整条复刻失败。此时还没预留视频积分(CREATED 阶段不预留),不扣费。
|
||
"""
|
||
from apps.ai.video_digest import VideoDigestError, digest_asset_video
|
||
|
||
if not is_video_replace_task(task) or task.status != AITask.Status.CREATED:
|
||
return
|
||
payload = task.request_payload or {}
|
||
if not payload.get("digest_pending"):
|
||
return
|
||
|
||
try:
|
||
asset = _load_replace_asset(
|
||
task.team, uuid.UUID(str(payload.get("digest_source_asset_id") or "")), "参考视频"
|
||
)
|
||
except (ValueError, TypeError):
|
||
_fail_reviewing_task(task, "参考视频已失效,请重新上传", error_code="asset_unavailable")
|
||
return
|
||
|
||
try:
|
||
digest, meta = digest_asset_video(asset=asset, task=task)
|
||
has_triview = any(
|
||
(ref or {}).get("label") == TRIVIEW_LABEL for ref in (payload.get("references") or [])
|
||
)
|
||
prompt = build_product_replace_prompt(
|
||
digest,
|
||
str(payload.get("subject_name") or ""),
|
||
has_triview=has_triview,
|
||
source_seconds=float(payload.get("digest_source_duration") or 0),
|
||
output_seconds=int(payload.get("duration") or 0),
|
||
)
|
||
except (VideoDigestError, ValueError) as exc:
|
||
_fail_reviewing_task(task, str(exc), error_code="processing_failed")
|
||
return
|
||
except Exception as exc: # noqa: BLE001 — 拆解任何异常都要把任务收尾,别留 CREATED 僵尸
|
||
logger.exception("video replace digest failed for %s", task.id)
|
||
_fail_reviewing_task(task, f"参考视频拆解失败:{exc}", error_code="processing_failed")
|
||
return
|
||
|
||
with transaction.atomic():
|
||
locked = AITask.objects.select_for_update().get(id=task.id)
|
||
if locked.status != AITask.Status.CREATED:
|
||
return
|
||
next_payload = dict(locked.request_payload or {})
|
||
next_payload.update(meta)
|
||
next_payload["digest_pending"] = False
|
||
next_payload["digest_text"] = digest[:32000]
|
||
next_payload["prompt"] = prompt
|
||
locked.request_payload = next_payload
|
||
locked.save(update_fields=["request_payload", "updated_at"])
|
||
task = locked
|
||
|
||
# 拆完就地推进一次:商品图多半早已过审,能直接提交火山,省掉一轮 8s 轮询。
|
||
try:
|
||
task = advance_video_replace(task)
|
||
except Exception: # noqa: BLE001 — 推进失败交给轮询重试,任务还在 CREATED
|
||
logger.warning("video replace advance after digest failed for %s", task.id, exc_info=True)
|
||
task.refresh_from_db()
|
||
if task.status == AITask.Status.CREATED:
|
||
_enqueue_replace_review_poll(task)
|
||
|
||
|
||
def _legacy_replace_mode(prompt: str) -> str:
|
||
return "character" if prompt.startswith("[视频复刻·角色]") else "product"
|
||
|
||
|
||
def _fail_reviewing_task(task, message: str, *, error_code: str = ""):
|
||
"""收尾一个还没提交火山的任务。
|
||
|
||
``_fail_pending_free_video`` 把错误码写死成 content_rejected —— 那是审核不通过的语义,
|
||
前端会渲染成「内容未通过生成审核」且不可重试。拆解失败/素材丢失不是那回事,
|
||
这里按实际原因改码,否则模型抖一下会被误报成合规问题,用户白白去换素材。
|
||
"""
|
||
from .free_video import _fail_pending_free_video
|
||
|
||
task = _fail_pending_free_video(task, message)
|
||
if error_code and task.status == AITask.Status.FAILED and task.error_code != error_code:
|
||
task.error_code = error_code
|
||
task.save(update_fields=["error_code", "updated_at"])
|
||
return task
|
||
|
||
|
||
def _enqueue_replace_review_poll(task):
|
||
try:
|
||
from .tasks import poll_free_video_task
|
||
|
||
poll_free_video_task.apply_async(args=[str(task.id), 0], countdown=8)
|
||
except Exception: # noqa: BLE001
|
||
logger.error("video replace review poll enqueue failed; relying on client polling", exc_info=True)
|
||
|
||
|
||
def _enqueue_replace_digest(task):
|
||
try:
|
||
from .tasks import run_video_replace_digest_task
|
||
|
||
run_video_replace_digest_task.delay(str(task.id))
|
||
except Exception: # noqa: BLE001
|
||
logger.error("video replace digest enqueue failed", exc_info=True)
|
||
|
||
|
||
def _create_reviewing_task(*, team, user, params: dict, digest_pending: bool = False):
|
||
"""还不能提交火山时:只建 CREATED 任务,不预留积分。
|
||
|
||
两种情况共用:①角色复刻的素材还在审核 ②商品复刻的参考视频还没拆解。
|
||
"""
|
||
model_name = str(params.get("model") or REPLACE_MODEL)
|
||
aspect_ratio = str(params.get("aspect_ratio") or "9:16")
|
||
resolution = str(params.get("resolution") or "720p")
|
||
try:
|
||
duration = int(params.get("duration") or 5)
|
||
except (TypeError, ValueError):
|
||
raise ValueError("时长参数无效")
|
||
if model_name not in FREE_VIDEO_MODELS:
|
||
raise ValueError("模型无效")
|
||
if aspect_ratio not in RATIOS:
|
||
raise ValueError("画面比例无效")
|
||
if resolution not in RESOLUTIONS:
|
||
raise ValueError("分辨率无效")
|
||
low, high = replace_duration_range()
|
||
if not low <= duration <= high:
|
||
raise ValueError(f"视频时长需在 {low}-{high} 秒之间")
|
||
|
||
model_config = (
|
||
ModelConfig.objects.select_related("provider")
|
||
.filter(name=model_name, capability=ModelConfig.Capability.VIDEO, status=ModelConfig.Status.ACTIVE)
|
||
.first()
|
||
)
|
||
if model_config is None:
|
||
raise ValueError("视频模型未配置,请联系管理员")
|
||
|
||
_reap_stale_free_video_tasks(team=team)
|
||
max_concurrent = int(getattr(settings, "FREE_VIDEO_MAX_CONCURRENT", 3))
|
||
in_flight = AITask.objects.filter(
|
||
team=team, task_type=AITask.Type.FREE_VIDEO, status__in=IN_FLIGHT_STATUSES
|
||
).count()
|
||
if in_flight >= max_concurrent:
|
||
raise ValueError(f"当前有 {in_flight} 个视频任务进行中(上限 {max_concurrent}),请等待完成后再提交")
|
||
|
||
references = params.get("references") or []
|
||
tokens, quote = quote_video_estimate(
|
||
model_config,
|
||
aspect_ratio=aspect_ratio,
|
||
resolution=resolution,
|
||
duration=duration,
|
||
references=references,
|
||
team=team,
|
||
)
|
||
reserve_amount = video_reserve_amount(quote.points)
|
||
account = CreditAccount.objects.filter(team=team).first()
|
||
available = (account.balance - account.reserved_balance) if account else Decimal("0")
|
||
if available < reserve_amount:
|
||
raise ValueError("团队余额不足,请充值后重试")
|
||
|
||
try:
|
||
seed = int(params.get("seed") if params.get("seed") is not None else -1)
|
||
except (TypeError, ValueError):
|
||
seed = -1
|
||
extra = params.get("extra_payload") if isinstance(params.get("extra_payload"), dict) else {}
|
||
request_payload = {
|
||
"feature": FEATURE,
|
||
"mode": "universal",
|
||
"model": model_name,
|
||
"endpoint": model_config.endpoint,
|
||
"prompt": params.get("prompt") or "",
|
||
"api_prompt": "",
|
||
"aspect_ratio": aspect_ratio,
|
||
"resolution": resolution,
|
||
"duration": duration,
|
||
"seed": seed,
|
||
"generate_audio": True,
|
||
"search_mode": "off",
|
||
"estimated_tokens": tokens,
|
||
"price_multiplier": quote.meta.get("price_multiplier", "1"),
|
||
"points_per_yuan_snapshot": quote.meta.get("rate", ""),
|
||
"references": references,
|
||
"model_routing_v1": True,
|
||
"review_pending": True,
|
||
"digest_pending": digest_pending,
|
||
}
|
||
for key, value in extra.items():
|
||
if key in request_payload or value in (None, ""):
|
||
continue
|
||
request_payload[key] = value
|
||
|
||
task = AITask.objects.create(
|
||
team=team,
|
||
created_by=user,
|
||
project=None,
|
||
task_type=AITask.Type.FREE_VIDEO,
|
||
status=AITask.Status.CREATED,
|
||
model_config=model_config,
|
||
idempotency_key=f"free_video:{team.id}:{uuid.uuid4()}",
|
||
request_payload=request_payload,
|
||
estimated_cost=quote.points,
|
||
base_cost=Decimal("0"),
|
||
)
|
||
if digest_pending:
|
||
_enqueue_replace_digest(task)
|
||
else:
|
||
_enqueue_replace_review_poll(task)
|
||
return task
|
||
|
||
|
||
def _replace_ref_assets(team, references: list) -> list[tuple[dict, Asset]]:
|
||
out = []
|
||
seen = set()
|
||
for ref in references or []:
|
||
raw_id = ref.get("asset_id")
|
||
if not raw_id:
|
||
continue
|
||
try:
|
||
parsed = uuid.UUID(str(raw_id))
|
||
except (TypeError, ValueError):
|
||
continue
|
||
if parsed in seen:
|
||
continue
|
||
seen.add(parsed)
|
||
asset = _load_replace_asset(team, parsed, ref.get("label") or "参考素材")
|
||
out.append((ref, asset))
|
||
return out
|
||
|
||
|
||
def _load_replace_asset(team, asset_id: uuid.UUID, label: str) -> Asset:
|
||
asset = Asset.objects.filter(id=asset_id, is_deleted=False, purged_at__isnull=True).first()
|
||
if asset is None:
|
||
raise ValueError(f"{label}不存在或已被删除")
|
||
if asset.team_id == team.id:
|
||
return asset
|
||
if Model.objects.filter(Q(is_official=True), Q(portrait_asset=asset) | Q(triview_asset=asset)).exists():
|
||
return asset
|
||
raise ValueError(f"{label}不存在或已被删除")
|
||
|
||
|
||
def _ensure_replace_refs_reviewed(team, references: list) -> str:
|
||
"""送审/轮询全部参考素材。返回 ready / pending / failed;审核未配置且无 remote_id 抛错。"""
|
||
from apps.assets import assets_client
|
||
from apps.assets.review import poll_asset_review, submit_asset_for_review
|
||
|
||
pairs = _replace_ref_assets(team, references)
|
||
if not pairs:
|
||
raise ValueError("请先上传参考视频")
|
||
states = []
|
||
for ref, asset in pairs:
|
||
label = ref.get("label") or asset.name or "参考素材"
|
||
if asset.review_status == "active" and asset.review_remote_id:
|
||
states.append("ready")
|
||
continue
|
||
if not assets_client.is_enabled():
|
||
raise ValueError(REVIEW_UNAVAILABLE)
|
||
if asset.review_status == "processing" and asset.review_remote_id:
|
||
poll_asset_review(asset)
|
||
asset.refresh_from_db(fields=["review_status", "review_remote_id", "review_error"])
|
||
elif asset.review_status != "active" or not asset.review_remote_id:
|
||
ok = submit_asset_for_review(asset, force=True)
|
||
asset.refresh_from_db(fields=["review_status", "review_remote_id", "review_error"])
|
||
if not ok and not asset.review_remote_id:
|
||
raise ValueError(f"「{label}」{REVIEW_SUBMIT_FAILED}")
|
||
if asset.review_status == "processing" and asset.review_remote_id:
|
||
poll_asset_review(asset)
|
||
asset.refresh_from_db(fields=["review_status", "review_remote_id", "review_error"])
|
||
if asset.review_status == "active" and asset.review_remote_id:
|
||
states.append("ready")
|
||
elif asset.review_status == "failed":
|
||
states.append("failed")
|
||
else:
|
||
states.append("pending")
|
||
if any(state == "failed" for state in states):
|
||
return "failed"
|
||
if all(state == "ready" for state in states):
|
||
return "ready"
|
||
return "pending"
|
||
|
||
|
||
def _refresh_replace_refs(team, references: list) -> list:
|
||
"""同一团队走 source=asset;官方跨团队素材把过审 id 写成 resolved_url=asset://。"""
|
||
out = []
|
||
for ref in references or []:
|
||
item = dict(ref)
|
||
raw_id = item.get("asset_id")
|
||
if not raw_id:
|
||
out.append(item)
|
||
continue
|
||
try:
|
||
parsed = uuid.UUID(str(raw_id))
|
||
except (TypeError, ValueError):
|
||
out.append(item)
|
||
continue
|
||
try:
|
||
asset = _load_replace_asset(team, parsed, item.get("label") or "参考素材")
|
||
except ValueError:
|
||
out.append(item)
|
||
continue
|
||
if asset.team_id == team.id:
|
||
item["source"] = "asset"
|
||
item.pop("resolved_url", None)
|
||
elif asset.review_remote_id:
|
||
item["source"] = "upload"
|
||
item["resolved_url"] = f"asset://{asset.review_remote_id}"
|
||
out.append(item)
|
||
return out
|
||
|
||
|
||
def _assert_replace_refs_ready(references: list) -> None:
|
||
missing = []
|
||
for ref in references or []:
|
||
raw_id = ref.get("asset_id")
|
||
if not raw_id:
|
||
continue
|
||
try:
|
||
parsed = uuid.UUID(str(raw_id))
|
||
except (TypeError, ValueError):
|
||
continue
|
||
asset = Asset.objects.filter(id=parsed).first()
|
||
if asset is None or not asset.review_remote_id or asset.review_status != "active":
|
||
missing.append(ref.get("label") or "参考素材")
|
||
if missing:
|
||
raise ValueError("素材尚未完成合规审核,请稍后再试")
|
||
|
||
|
||
def _optional_uuid(value, label: str):
|
||
text = str(value or "").strip()
|
||
if not text:
|
||
return None
|
||
try:
|
||
return uuid.UUID(text)
|
||
except (TypeError, ValueError) as exc:
|
||
raise ValueError(f"{label}无效") from exc
|
||
|
||
|
||
def _uuid_list(value, label: str) -> list:
|
||
if value in (None, ""):
|
||
return []
|
||
if not isinstance(value, (list, tuple)):
|
||
raise ValueError(f"{label}格式无效")
|
||
if len(value) > MAX_IMAGES:
|
||
raise ValueError(f"{label}最多 {MAX_IMAGES} 张")
|
||
seen = set()
|
||
out = []
|
||
for item in value:
|
||
parsed = _optional_uuid(item, label)
|
||
if parsed is None or parsed in seen:
|
||
continue
|
||
seen.add(parsed)
|
||
out.append(parsed)
|
||
return out
|
||
|
||
|
||
def _team_asset(team, asset_id, *, kind: str, label: str) -> Asset:
|
||
parsed = _optional_uuid(asset_id, label)
|
||
if parsed is None:
|
||
raise ValueError(f"请先上传{label}")
|
||
asset = Asset.objects.filter(id=parsed, team=team, is_deleted=False, purged_at__isnull=True).first()
|
||
if asset is None:
|
||
raise ValueError(f"{label}不存在或已被删除")
|
||
if asset.asset_type != kind:
|
||
raise ValueError(f"{label}类型不正确")
|
||
return asset
|
||
|
||
|
||
def _asset_duration_seconds(asset: Asset) -> float:
|
||
primary = asset.files.filter(is_primary=True).first() or asset.files.first()
|
||
if primary is None or not primary.duration_ms:
|
||
return 0.0
|
||
return primary.duration_ms / 1000.0
|
||
|
||
|
||
def replace_duration_range() -> tuple[int, int]:
|
||
"""复刻模型(Seedance 2.5)支持的出片时长区间。模型没配好时回落 4–15,不至于炸。"""
|
||
model = (
|
||
ModelConfig.objects.filter(
|
||
name=REPLACE_MODEL, capability=ModelConfig.Capability.VIDEO, status=ModelConfig.Status.ACTIVE
|
||
).first()
|
||
)
|
||
return model_duration_range(model)
|
||
|
||
|
||
def _output_duration(requested, video_seconds: float) -> int:
|
||
low, high = replace_duration_range()
|
||
try:
|
||
value = int(requested) if requested not in (None, "") else 0
|
||
except (TypeError, ValueError):
|
||
value = 0
|
||
if value:
|
||
return min(high, max(low, value))
|
||
if video_seconds:
|
||
return min(high, max(low, int(round(video_seconds))))
|
||
return high
|
||
|
||
|
||
def _owned_ref(asset: Asset, *, kind: str, role: str, label: str) -> dict:
|
||
from .services import _asset_preview_url
|
||
|
||
url = _asset_preview_url(asset)
|
||
if not url:
|
||
raise ValueError(f"「{label}」没有可用文件")
|
||
ref = {
|
||
"url": url,
|
||
"type": kind,
|
||
"role": role,
|
||
"label": label,
|
||
"source": "asset",
|
||
"asset_id": str(asset.id),
|
||
}
|
||
seconds = _asset_duration_seconds(asset)
|
||
if seconds:
|
||
ref["duration"] = seconds
|
||
return ref
|
||
|
||
|
||
def _library_image_ref(asset: Asset, *, team, label: str) -> dict:
|
||
from .services import _asset_preview_url, _seedance_ref_url
|
||
|
||
if asset.team_id == team.id and not asset.is_deleted:
|
||
return {
|
||
"url": _asset_preview_url(asset) or "",
|
||
"type": "image",
|
||
"role": "reference_image",
|
||
"label": label,
|
||
"source": "asset",
|
||
"asset_id": str(asset.id),
|
||
}
|
||
raw = _asset_preview_url(asset)
|
||
url = _seedance_ref_url(raw, asset.review_status, asset.review_remote_id)
|
||
if not url:
|
||
raise ValueError(f"「{label}」没有可用文件")
|
||
return {
|
||
"url": url,
|
||
"type": "image",
|
||
"role": "reference_image",
|
||
"label": label,
|
||
"source": "upload",
|
||
"asset_id": str(asset.id),
|
||
}
|
||
|
||
|
||
def _product_triview_asset(team, product_id: uuid.UUID):
|
||
"""商品库里这个商品的三视图(白底多角度单张 16:9)。没有就返回 None,不挡生成。"""
|
||
return (
|
||
Asset.objects.filter(
|
||
team=team,
|
||
is_deleted=False,
|
||
purged_at__isnull=True,
|
||
metadata__product_id=str(product_id),
|
||
metadata__view="three_view",
|
||
)
|
||
.order_by("-created_at")
|
||
.first()
|
||
)
|
||
|
||
|
||
def _product_library_refs(team, product_id: uuid.UUID) -> tuple[str, list, str]:
|
||
product = (
|
||
Product.objects.filter(id=product_id, team=team, purged_at__isnull=True, status=Product.Status.ACTIVE)
|
||
.select_related("cover_asset")
|
||
.prefetch_related("images__asset")
|
||
.first()
|
||
)
|
||
if product is None:
|
||
raise ValueError("商品不存在或已被删除")
|
||
triview = _product_triview_asset(team, product_id)
|
||
# 三视图也占一张参考图的名额,商品图要给它让位,否则火山那边超 9 张直接拒。
|
||
image_budget = MAX_IMAGES - 1 if triview is not None else MAX_IMAGES
|
||
assets = []
|
||
seen = {triview.id} if triview is not None else set()
|
||
for image in product.images.all():
|
||
asset = image.asset
|
||
if asset is None or asset.id in seen or asset.is_deleted:
|
||
continue
|
||
seen.add(asset.id)
|
||
assets.append(asset)
|
||
if len(assets) >= image_budget:
|
||
break
|
||
if not assets and product.cover_asset_id and product.cover_asset_id not in seen and not product.cover_asset.is_deleted:
|
||
assets.append(product.cover_asset)
|
||
if not assets:
|
||
# 只有三视图、没有任何实拍图时,把三视图顶上来当 @目标商品。
|
||
# 否则提示词里的 @目标商品 找不到同名 label,火山拿到的是一句没有指代的字面量。
|
||
if triview is None:
|
||
raise ValueError("这个商品还没有可用图片")
|
||
assets = [triview]
|
||
triview = None
|
||
refs = [_library_image_ref(asset, team=team, label="目标商品" if index == 0 else f"目标商品{index + 1}") for index, asset in enumerate(assets)]
|
||
if triview is not None:
|
||
# 放最后:@目标商品 仍指向第一张实拍图,三视图作为「各面长什么样」的补充证据。
|
||
refs.append(_library_image_ref(triview, team=team, label=TRIVIEW_LABEL))
|
||
return product.title, refs, "library"
|
||
|
||
|
||
def _character_library_refs(team, model_id: uuid.UUID) -> tuple[str, list, str]:
|
||
model = (
|
||
Model.objects.filter(Q(team=team) | Q(is_official=True), id=model_id, is_deleted=False, purged_at__isnull=True)
|
||
.select_related("portrait_asset", "triview_asset")
|
||
.first()
|
||
)
|
||
if model is None:
|
||
raise ValueError("角色不存在或已被删除")
|
||
assets = []
|
||
seen = set()
|
||
for asset in (model.portrait_asset, model.triview_asset):
|
||
if asset is None or asset.id in seen or asset.is_deleted:
|
||
continue
|
||
seen.add(asset.id)
|
||
assets.append(asset)
|
||
if len(assets) >= MAX_IMAGES:
|
||
break
|
||
if not assets:
|
||
raise ValueError("这个角色还没有可用图片")
|
||
labels = ["目标角色", "目标角色三视图"]
|
||
refs = [_library_image_ref(asset, team=team, label=labels[index] if index < len(labels) else f"目标角色{index + 1}") for index, asset in enumerate(assets)]
|
||
return model.name, refs, "library"
|
||
|
||
|
||
def _temporary_image_refs(team, image_ids: list, *, noun: str) -> tuple[str, list, str]:
|
||
refs = []
|
||
for index, asset_id in enumerate(image_ids):
|
||
asset = _team_asset(team, asset_id, kind=Asset.Type.IMAGE, label=f"{noun}参考图")
|
||
label = "目标商品" if noun == "商品" else "目标角色"
|
||
if index > 0:
|
||
label = f"{label}{index + 1}"
|
||
refs.append(_owned_ref(asset, kind="image", role="reference_image", label=label))
|
||
fallback = "临时商品素材" if noun == "商品" else "临时角色素材"
|
||
name = Asset.objects.filter(id=image_ids[0]).values_list("name", flat=True).first() or fallback
|
||
subject = name.rsplit(".", 1)[0] if name else fallback
|
||
if len(refs) > 1:
|
||
subject = f"{subject}({len(refs)}张参考图)"
|
||
return subject, refs, "temporary"
|