完善视频复刻

This commit is contained in:
Azmat@qq.com
2026-08-27 15:56:44 +08:00
parent 3fbc2f2dab
commit d2b786a3cf
10 changed files with 722 additions and 164 deletions
+140 -10
View File
@@ -47,6 +47,7 @@ RATIOS = {"21:9", "16:9", "4:3", "1:1", "3:4", "9:16"}
RESOLUTIONS = {"480p", "720p", "1080p", "4k"}
MODES = {"universal", "keyframe"}
IN_FLIGHT_STATUSES = (
AITask.Status.CREATED, # 视频复刻审核中也占并发,避免连点刷出一堆待审任务
AITask.Status.RESERVED,
AITask.Status.SUBMITTED,
AITask.Status.POLLING,
@@ -296,15 +297,16 @@ def build_content_items(*, team, prompt: str, mode: str, references: list) -> di
label_to_placeholder[label] = _placeholder_for(asset_type)
continue
# 直传素材(已上传 TOS 的直链)
# 直传素材(已上传 TOS 的直链)。resolved_url 可覆盖为 asset://(官方模特跨团队等)
push_url = str(ref.get("resolved_url") or url)
if ref_type == "image":
# 参考图模式下所有图 role 必须 reference_image;keyframe 用 first_frame/last_frame
effective_role = "reference_image" if mode == "universal" else (role or "first_frame")
asset_type = _push("image", url, effective_role)
asset_type = _push("image", push_url, effective_role)
elif ref_type == "video":
asset_type = _push("video", url, role or "reference_video", duration)
asset_type = _push("video", push_url, role or "reference_video", duration)
elif ref_type == "audio":
asset_type = _push("audio", url, role or "reference_audio", duration)
asset_type = _push("audio", push_url, role or "reference_audio", duration)
else:
logger.warning("unknown ref_type=%s url=%s label=%s, skipped", ref_type, url, label)
continue
@@ -366,6 +368,11 @@ def _reap_stale_free_video_tasks(*, team) -> None:
video_policy = load_model_routing_policy().video
buckets = [
(
[AITask.Status.CREATED],
{"updated_at__lt": now - timedelta(minutes=16)},
"素材审核超时(自动回收)",
),
(
[AITask.Status.RESERVED],
{"updated_at__lt": now - timedelta(seconds=video_policy.submit_total_timeout)},
@@ -538,7 +545,133 @@ def submit_free_video(*, team, user, params: dict) -> AITask:
task.status = AITask.Status.RESERVED
task.save(update_fields=["status", "updated_at"])
# 火山调用在事务外(不持锁调外网)
return _dispatch_free_video_provider(
task=task,
built=built,
model_config=model_config,
aspect_ratio=aspect_ratio,
duration=duration,
resolution=resolution,
generate_audio=generate_audio,
seed=seed,
search_mode=search_mode,
feature=feature,
mode=mode,
poll_countdown=30,
)
def start_pending_free_video(task: AITask) -> AITask:
"""CREATED 任务审核已过后:预留积分 + 调火山。并发 poll 用行锁认领,失败不扣费。"""
if task.status != AITask.Status.CREATED:
return task
payload = task.request_payload or {}
prompt = str(payload.get("prompt") or "").strip()
mode = str(payload.get("mode") or "universal")
aspect_ratio = str(payload.get("aspect_ratio") or "16:9")
resolution = str(payload.get("resolution") or "720p")
generate_audio = bool(payload.get("generate_audio", True))
search_mode = str(payload.get("search_mode") or "off")
feature = str(payload.get("feature") or "free_video")
try:
duration = int(payload.get("duration") or 5)
except (TypeError, ValueError):
duration = 5
try:
seed = int(payload.get("seed") if payload.get("seed") is not None else -1)
except (TypeError, ValueError):
seed = -1
references = payload.get("references") or []
try:
built = build_content_items(team=task.team, prompt=prompt, mode=mode, references=references)
except ValueError as exc:
message = str(exc)
if "正在审核" in message or "已提交审核" in message or "尚未完成合规审核" in message:
return task
return _fail_pending_free_video(task, message)
tokens, quote = quote_video_estimate(
task.model_config,
aspect_ratio=aspect_ratio,
resolution=resolution,
duration=duration,
references=built["snapshots"],
team=task.team,
)
reserve_amount = video_reserve_amount(quote.points)
with transaction.atomic():
locked = (
AITask.objects.select_for_update()
.select_related("model_config", "model_config__provider", "team", "created_by")
.get(id=task.id)
)
if locked.status != AITask.Status.CREATED:
return locked
try:
reserve_credit(team=locked.team, user=locked.created_by, task=locked, amount=reserve_amount)
except ValueError as exc:
message = "团队余额不足,请充值后重试" if "insufficient credit" in str(exc) else str(exc)
locked.status = AITask.Status.FAILED
locked.error_code = "user_credit_insufficient" if "余额不足" in message else "invalid_input"
locked.error_message = message[:2000]
locked.completed_at = timezone.now()
locked.save(update_fields=["status", "error_code", "error_message", "completed_at", "updated_at"])
return locked
next_payload = dict(locked.request_payload or {})
next_payload["api_prompt"] = built["api_prompt"]
next_payload["references"] = built["snapshots"]
next_payload["estimated_tokens"] = tokens
next_payload["review_pending"] = False
locked.request_payload = next_payload
locked.estimated_cost = quote.points
locked.status = AITask.Status.RESERVED
locked.save(update_fields=["request_payload", "estimated_cost", "status", "updated_at"])
return _dispatch_free_video_provider(
task=locked,
built=built,
model_config=locked.model_config,
aspect_ratio=aspect_ratio,
duration=duration,
resolution=resolution,
generate_audio=generate_audio,
seed=seed,
search_mode=search_mode,
feature=feature,
mode=mode,
poll_countdown=30,
)
def _fail_pending_free_video(task: AITask, message: str) -> AITask:
if task.status not in (AITask.Status.CREATED, AITask.Status.RESERVED):
return task
task.status = AITask.Status.FAILED
task.error_code = "content_rejected"
task.error_message = message[:2000]
task.completed_at = timezone.now()
task.save(update_fields=["status", "error_code", "error_message", "completed_at", "updated_at"])
_notify_failure(task, raw=message, hint=message)
return task
def _dispatch_free_video_provider(
*,
task,
built,
model_config,
aspect_ratio,
duration,
resolution,
generate_audio,
seed,
search_mode,
feature,
mode,
poll_countdown=30,
):
"""RESERVED 任务调火山创建。失败退费。"""
try:
from .services import execute_routed_video_submit
@@ -558,8 +691,6 @@ def submit_free_video(*, team, user, params: dict) -> AITask:
request_summary={"feature": feature, "mode": mode},
)
response, provider_task_id = routed.value
# execute_model_call 以原子 F 表达式累计实际尝试的平台成本;刷新内存对象,
# 保证本方法返回值与数据库中的 AITask.base_cost 完全一致。
task.refresh_from_db(fields=["base_cost"])
task.provider_task_id = provider_task_id
task.response_payload = response
@@ -596,7 +727,7 @@ def submit_free_video(*, team, user, params: dict) -> AITask:
from apps.assets.free_asset_state import mark_remote_asset_unavailable
mark_remote_asset_unavailable(
team=team,
team=task.team,
local_asset_id=target["local_asset_id"],
remote_asset_id=target["remote_asset_id"],
)
@@ -616,11 +747,10 @@ def submit_free_video(*, team, user, params: dict) -> AITask:
logger.warning("free video create failed: %s", exc)
return task
# worker 兜底轮询(自重排);派发失败仅 log,前端主动 poll 仍能收尾
try:
from .tasks import poll_free_video_task
poll_free_video_task.apply_async(args=[str(task.id), 0], countdown=30)
poll_free_video_task.apply_async(args=[str(task.id), 0], countdown=poll_countdown)
except Exception: # noqa: BLE001
logger.error("poll_free_video_task enqueue failed; relying on client polling", exc_info=True)